diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java index cb90fa4345..e90ae5c383 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java @@ -26,8 +26,6 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.id.EdgeId; -import org.thingsboard.server.common.data.id.RuleChainId; -import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.EdgeRpcServiceGrpc; import org.thingsboard.server.gen.edge.RequestMsg; import org.thingsboard.server.gen.edge.ResponseMsg; @@ -48,7 +46,7 @@ import java.util.concurrent.Executors; public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase { private final Map sessions = new ConcurrentHashMap<>(); - private static final ObjectMapper objectMapper = new ObjectMapper(); + private static final ObjectMapper mapper = new ObjectMapper(); @Value("${edges.rpc.port}") private int rpcPort; @@ -102,7 +100,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase { @Override public StreamObserver handleMsgs(StreamObserver outputStream) { - return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, objectMapper).getInputStream(); + return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, mapper).getInputStream(); } private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { 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 341142e34d..6d2637840a 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 @@ -18,7 +18,6 @@ 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.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -55,7 +54,6 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; -import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.page.TimePageData; import org.thingsboard.server.common.data.page.TimePageLink; @@ -70,14 +68,18 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; +import org.thingsboard.server.gen.edge.AttributesRequestMsg; 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.DeviceCredentialsRequestMsg; +import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; 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.EntityDataProto; import org.thingsboard.server.gen.edge.EntityUpdateMsg; +import org.thingsboard.server.gen.edge.RelationRequestMsg; import org.thingsboard.server.gen.edge.RequestMsg; import org.thingsboard.server.gen.edge.RequestMsgType; import org.thingsboard.server.gen.edge.ResponseMsg; @@ -86,6 +88,7 @@ import org.thingsboard.server.gen.edge.RuleChainMetadataUpdateMsg; 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.gen.edge.UserCredentialsRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.edge.EdgeContextComponent; @@ -115,7 +118,7 @@ public final class EdgeGrpcSession implements Closeable { private final UUID sessionId; private final BiConsumer sessionOpenListener; private final Consumer sessionCloseListener; - private final ObjectMapper objectMapper; + private final ObjectMapper mapper; private EdgeContextComponent ctx; private Edge edge; @@ -124,13 +127,13 @@ public final class EdgeGrpcSession implements Closeable { private boolean connected; EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, BiConsumer sessionOpenListener, - Consumer sessionCloseListener, ObjectMapper objectMapper) { + Consumer sessionCloseListener, ObjectMapper mapper) { this.sessionId = UUID.randomUUID(); this.ctx = ctx; this.outputStream = outputStream; this.sessionOpenListener = sessionOpenListener; this.sessionCloseListener = sessionCloseListener; - this.objectMapper = objectMapper; + this.mapper = mapper; initInputStream(); } @@ -147,7 +150,7 @@ public final class EdgeGrpcSession implements Closeable { outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); } if (ConnectResponseCode.ACCEPTED == responseMsg.getResponseCode()) { - ctx.getSyncEdgeService().sync(ctx, edge, outputStream); + ctx.getSyncEdgeService().sync(edge, outputStream); } } if (connected) { @@ -189,9 +192,6 @@ public final class EdgeGrpcSession implements Closeable { processTelemetryMessage(edgeEvent); } else { processEntityCRUDMessage(edgeEvent, msgType); - if (ENTITY_CREATED_RPC_MESSAGE.equals(msgType)) { - pushEntityAttributesToEdge(edgeEvent); - } } } catch (Exception e) { log.error("Exception during processing records from queue", e); @@ -239,58 +239,6 @@ public final class EdgeGrpcSession implements Closeable { ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); } - private void pushEntityAttributesToEdge(EdgeEvent edgeEvent) throws IOException { - EntityId entityId = null; - switch (edgeEvent.getEdgeEventType()) { - case EDGE: - entityId = edge.getId(); - break; - case DEVICE: - entityId = new DeviceId(edgeEvent.getEntityId()); - break; - case ASSET: - entityId = new AssetId(edgeEvent.getEntityId()); - break; - case ENTITY_VIEW: - entityId = new EntityViewId(edgeEvent.getEntityId()); - break; - case DASHBOARD: - entityId = new DashboardId(edgeEvent.getEntityId()); - break; - } - if (entityId != null) { - final EntityId finalEntityId = entityId; - ListenableFuture> ssAttrFuture = ctx.getAttributesService().findAll(edge.getTenantId(), entityId, DataConstants.SERVER_SCOPE); - Futures.transform(ssAttrFuture, ssAttributes -> { - if (ssAttributes != null && !ssAttributes.isEmpty()) { - try { - ObjectNode entityNode = objectMapper.createObjectNode(); - for (AttributeKvEntry attr : ssAttributes) { - if (attr.getDataType() == DataType.BOOLEAN && attr.getBooleanValue().isPresent()) { - entityNode.put(attr.getKey(), attr.getBooleanValue().get()); - } else if (attr.getDataType() == DataType.DOUBLE && attr.getDoubleValue().isPresent()) { - entityNode.put(attr.getKey(), attr.getDoubleValue().get()); - } else if (attr.getDataType() == DataType.LONG && attr.getLongValue().isPresent()) { - entityNode.put(attr.getKey(), attr.getLongValue().get()); - } else { - entityNode.put(attr.getKey(), attr.getValueAsString()); - } - } - log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", finalEntityId, entityNode); - DownlinkMsg value = constructEntityDataProtoMsg(finalEntityId, ActionType.ATTRIBUTES_UPDATED, JsonUtils.parse(objectMapper.writeValueAsString(entityNode))); - outputStream.onNext(ResponseMsg.newBuilder() - .setDownlinkMsg(value).build()); - } catch (Exception e) { - log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e); - } - } - return null; - }, MoreExecutors.directExecutor()); - ListenableFuture> shAttrFuture = ctx.getAttributesService().findAll(edge.getTenantId(), entityId, DataConstants.SHARED_SCOPE); - ListenableFuture> clAttrFuture = ctx.getAttributesService().findAll(edge.getTenantId(), entityId, DataConstants.CLIENT_SCOPE); - } - } - private void processTelemetryMessage(EdgeEvent edgeEvent) throws IOException { log.trace("Executing processTelemetryMessage, edgeEvent [{}]", edgeEvent); EntityId entityId = null; @@ -311,7 +259,7 @@ public final class EdgeGrpcSession implements Closeable { DownlinkMsg downlinkMsg; try { ActionType actionType = ActionType.valueOf(edgeEvent.getEdgeEventAction()); - downlinkMsg = constructEntityDataProtoMsg(entityId, actionType, JsonUtils.parse(objectMapper.writeValueAsString(edgeEvent.getEntityBody()))); + downlinkMsg = constructEntityDataProtoMsg(entityId, actionType, JsonUtils.parse(mapper.writeValueAsString(edgeEvent.getEntityBody()))); outputStream.onNext(ResponseMsg.newBuilder() .setDownlinkMsg(downlinkMsg) .build()); @@ -617,7 +565,7 @@ public final class EdgeGrpcSession implements Closeable { } private void processRelationCRUD(EdgeEvent edgeEvent, UpdateMsgType msgType) { - EntityRelation entityRelation = objectMapper.convertValue(edgeEvent.getEntityBody(), EntityRelation.class); + EntityRelation entityRelation = mapper.convertValue(edgeEvent.getEntityBody(), EntityRelation.class); EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() .setRelationUpdateMsg(ctx.getRelationUpdateMsgConstructor().constructRelationUpdatedMsg(msgType, entityRelation)) .build(); @@ -666,7 +614,7 @@ public final class EdgeGrpcSession implements Closeable { return UpdateMsgType.ALARM_CLEAR_RPC_MESSAGE; case ATTRIBUTES_UPDATED: case ATTRIBUTES_DELETED: - case TIMESERIES_DELETED: + case TIMESERIES_UPDATED: return null; default: throw new RuntimeException("Unsupported actionType [" + actionType + "]"); @@ -702,11 +650,17 @@ public final class EdgeGrpcSession implements Closeable { } } } + if (uplinkMsg.getDeviceUpdateMsgList() != null && !uplinkMsg.getDeviceUpdateMsgList().isEmpty()) { for (DeviceUpdateMsg deviceUpdateMsg : uplinkMsg.getDeviceUpdateMsgList()) { onDeviceUpdate(deviceUpdateMsg); } } + if (uplinkMsg.getDeviceCredentialsUpdateMsgList() != null && !uplinkMsg.getDeviceCredentialsUpdateMsgList().isEmpty()) { + for (DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg : uplinkMsg.getDeviceCredentialsUpdateMsgList()) { + onDeviceCredentialsUpdate(deviceCredentialsUpdateMsg); + } + } if (uplinkMsg.getAlarmUpdateMsgList() != null && !uplinkMsg.getAlarmUpdateMsgList().isEmpty()) { for (AlarmUpdateMsg alarmUpdateMsg : uplinkMsg.getAlarmUpdateMsgList()) { onAlarmUpdate(alarmUpdateMsg); @@ -714,7 +668,27 @@ public final class EdgeGrpcSession implements Closeable { } if (uplinkMsg.getRuleChainMetadataRequestMsgList() != null && !uplinkMsg.getRuleChainMetadataRequestMsgList().isEmpty()) { for (RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg : uplinkMsg.getRuleChainMetadataRequestMsgList()) { - ctx.getSyncEdgeService().syncRuleChainMetadata(edge, ruleChainMetadataRequestMsg, outputStream); + ctx.getSyncEdgeService().processRuleChainMetadata(edge, ruleChainMetadataRequestMsg, outputStream); + } + } + if (uplinkMsg.getAttributesRequestMsgList() != null && !uplinkMsg.getAttributesRequestMsgList().isEmpty()) { + for (AttributesRequestMsg attributesRequestMsg : uplinkMsg.getAttributesRequestMsgList()) { + ctx.getSyncEdgeService().processAttributesRequestMsg(edge, attributesRequestMsg, outputStream); + } + } + if (uplinkMsg.getRelationRequestMsgList() != null && !uplinkMsg.getRelationRequestMsgList().isEmpty()) { + for (RelationRequestMsg relationRequestMsg : uplinkMsg.getRelationRequestMsgList()) { + ctx.getSyncEdgeService().processRelationRequestMsg(edge, relationRequestMsg, outputStream); + } + } + if (uplinkMsg.getUserCredentialsRequestMsgList() != null && !uplinkMsg.getUserCredentialsRequestMsgList().isEmpty()) { + for (UserCredentialsRequestMsg userCredentialsRequestMsg : uplinkMsg.getUserCredentialsRequestMsgList()) { + ctx.getSyncEdgeService().processUserCredentialsRequestMsg(edge, userCredentialsRequestMsg, outputStream); + } + } + if (uplinkMsg.getDeviceCredentialsRequestMsgList() != null && !uplinkMsg.getDeviceCredentialsRequestMsgList().isEmpty()) { + for (DeviceCredentialsRequestMsg deviceCredentialsRequestMsg : uplinkMsg.getDeviceCredentialsRequestMsgList()) { + ctx.getSyncEdgeService().processDeviceCredentialsRequestMsg(edge, deviceCredentialsRequestMsg, outputStream); } } } catch (Exception e) { @@ -850,21 +824,59 @@ public final class EdgeGrpcSession implements Closeable { device.setType(deviceUpdateMsg.getType()); device.setLabel(deviceUpdateMsg.getLabel()); device = ctx.getDeviceService().saveDevice(device); - updateDeviceCredentials(deviceUpdateMsg, device); + + requestDeviceCredentialsFromEdge(device); + } + + private void onDeviceCredentialsUpdate(DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg) { + log.debug("Executing onDeviceCredentialsUpdate, deviceCredentialsUpdateMsg [{}]", deviceCredentialsUpdateMsg); + DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsUpdateMsg.getDeviceIdMSB(), deviceCredentialsUpdateMsg.getDeviceIdLSB())); + ListenableFuture deviceFuture = ctx.getDeviceService().findDeviceByIdAsync(edge.getTenantId(), deviceId); + + Futures.addCallback(deviceFuture, new FutureCallback() { + @Override + public void onSuccess(@Nullable Device device) { + if (device != null) { + log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", + device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue()); + try { + DeviceCredentials deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(edge.getTenantId(), device.getId()); + deviceCredentials.setCredentialsType(DeviceCredentialsType.valueOf(deviceCredentialsUpdateMsg.getCredentialsType())); + deviceCredentials.setCredentialsId(deviceCredentialsUpdateMsg.getCredentialsId()); + deviceCredentials.setCredentialsValue(deviceCredentialsUpdateMsg.getCredentialsValue()); + ctx.getDeviceCredentialsService().updateDeviceCredentials(edge.getTenantId(), deviceCredentials); + } catch (Exception e) { + log.error("Can't update device credentials for device [{}], deviceCredentialsUpdateMsg [{}]", device.getName(), deviceCredentialsUpdateMsg, e); + } + log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", + device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue()); + } + } + + @Override + public void onFailure(Throwable t) { + log.error("Can't update device credentials for deviceCredentialsUpdateMsg [{}]", deviceCredentialsUpdateMsg, t); + } + }, ctx.getDbCallbackExecutor()); } - private void updateDeviceCredentials(DeviceUpdateMsg deviceUpdateMsg, Device device) { - log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", - device.getName(), deviceUpdateMsg.getCredentialsId(), deviceUpdateMsg.getCredentialsValue()); + private void requestDeviceCredentialsFromEdge(Device device) { + log.debug("Executing requestDeviceCredentialsFromEdge device [{}]", device); - DeviceCredentials deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(edge.getTenantId(), device.getId()); - deviceCredentials.setCredentialsType(DeviceCredentialsType.valueOf(deviceUpdateMsg.getCredentialsType())); - deviceCredentials.setCredentialsId(deviceUpdateMsg.getCredentialsId()); - deviceCredentials.setCredentialsValue(deviceUpdateMsg.getCredentialsValue()); - ctx.getDeviceCredentialsService().updateDeviceCredentials(edge.getTenantId(), deviceCredentials); - log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", - device.getName(), deviceUpdateMsg.getCredentialsId(), deviceUpdateMsg.getCredentialsValue()); + DownlinkMsg downlinkMsg = constructDeviceCredentialsRequestMsg(device.getId()); + outputStream.onNext(ResponseMsg.newBuilder() + .setDownlinkMsg(downlinkMsg) + .build()); + } + private DownlinkMsg constructDeviceCredentialsRequestMsg(DeviceId deviceId) { + DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = DeviceCredentialsRequestMsg.newBuilder() + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .build(); + DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() + .addAllDeviceCredentialsRequestMsg(Collections.singletonList(deviceCredentialsRequestMsg)); + return builder.build(); } private Device createDevice(DeviceUpdateMsg deviceUpdateMsg) { @@ -884,7 +896,8 @@ public final class EdgeGrpcSession implements Closeable { createRelationFromEdge(device.getId()); ctx.getRelationService().saveRelationAsync(TenantId.SYS_TENANT_ID, new EntityRelation(edge.getId(), device.getId(), "Created")); ctx.getDeviceStateService().onDeviceAdded(device); - updateDeviceCredentials(deviceUpdateMsg, device); + + requestDeviceCredentialsFromEdge(device); } finally { deviceCreationLock.unlock(); } @@ -925,7 +938,7 @@ public final class EdgeGrpcSession implements Closeable { existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); } existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); - existentAlarm.setDetails(objectMapper.readTree(alarmUpdateMsg.getDetails())); + existentAlarm.setDetails(mapper.readTree(alarmUpdateMsg.getDetails())); ctx.getAlarmService().createOrUpdateAlarm(existentAlarm); break; case ALARM_ACK_RPC_MESSAGE: @@ -935,7 +948,7 @@ public final class EdgeGrpcSession implements Closeable { break; case ALARM_CLEAR_RPC_MESSAGE: if (existentAlarm != null) { - ctx.getAlarmService().clearAlarm(edge.getTenantId(), existentAlarm.getId(), objectMapper.readTree(alarmUpdateMsg.getDetails()), alarmUpdateMsg.getAckTs()); + ctx.getAlarmService().clearAlarm(edge.getTenantId(), existentAlarm.getId(), mapper.readTree(alarmUpdateMsg.getDetails()), alarmUpdateMsg.getAckTs()); } break; case ENTITY_DELETED_RPC_MESSAGE: diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardUpdateMsgConstructor.java index 6bd24f123b..6c52464326 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardUpdateMsgConstructor.java @@ -16,11 +16,9 @@ package org.thingsboard.server.service.edge.rpc.constructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.id.DashboardId; -import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.dao.util.mapping.JacksonUtil; import org.thingsboard.server.gen.edge.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.UpdateMsgType; @@ -29,11 +27,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Slf4j public class DashboardUpdateMsgConstructor { - @Autowired - private DashboardService dashboardService; - public DashboardUpdateMsg constructDashboardUpdatedMsg(UpdateMsgType msgType, Dashboard dashboard) { - dashboard = dashboardService.findDashboardById(dashboard.getTenantId(), dashboard.getId()); DashboardUpdateMsg.Builder builder = DashboardUpdateMsg.newBuilder() .setMsgType(msgType) .setIdMSB(dashboard.getId().getId().getMostSignificantBits()) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceUpdateMsgConstructor.java index 8d7fbc4c25..dcef0a57c6 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceUpdateMsgConstructor.java @@ -16,12 +16,11 @@ package org.thingsboard.server.service.edge.rpc.constructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.security.DeviceCredentials; -import org.thingsboard.server.dao.device.DeviceCredentialsService; +import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.UpdateMsgType; @@ -29,9 +28,6 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Slf4j public class DeviceUpdateMsgConstructor { - @Autowired - private DeviceCredentialsService deviceCredentialsService; - public DeviceUpdateMsg constructDeviceUpdatedMsg(UpdateMsgType msgType, Device device) { DeviceUpdateMsg.Builder builder = DeviceUpdateMsg.newBuilder() .setMsgType(msgType) @@ -42,20 +38,19 @@ public class DeviceUpdateMsgConstructor { if (device.getLabel() != null) { builder.setLabel(device.getLabel()); } - if (msgType.equals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE) || - msgType.equals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE) || - msgType.equals(UpdateMsgType.DEVICE_CONFLICT_RPC_MESSAGE)) { - DeviceCredentials deviceCredentials - = deviceCredentialsService.findDeviceCredentialsByDeviceId(device.getTenantId(), device.getId()); - if (deviceCredentials != null) { - if (deviceCredentials.getCredentialsType() != null) { - builder.setCredentialsType(deviceCredentials.getCredentialsType().name()) - .setCredentialsId(deviceCredentials.getCredentialsId()); - } - if (deviceCredentials.getCredentialsValue() != null) { - builder.setCredentialsValue(deviceCredentials.getCredentialsValue()); - } - } + return builder.build(); + } + + public DeviceCredentialsUpdateMsg constructDeviceCredentialsUpdatedMsg(DeviceCredentials deviceCredentials) { + DeviceCredentialsUpdateMsg.Builder builder = DeviceCredentialsUpdateMsg.newBuilder() + .setDeviceIdMSB(deviceCredentials.getDeviceId().getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceCredentials.getDeviceId().getId().getLeastSignificantBits()); + if (deviceCredentials.getCredentialsType() != null) { + builder.setCredentialsType(deviceCredentials.getCredentialsType().name()) + .setCredentialsId(deviceCredentials.getCredentialsId()); + } + if (deviceCredentials.getCredentialsValue() != null) { + builder.setCredentialsValue(deviceCredentials.getCredentialsValue()); } return builder.build(); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserUpdateMsgConstructor.java index b994d8fdc4..af42448e19 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserUpdateMsgConstructor.java @@ -16,31 +16,26 @@ package org.thingsboard.server.service.edge.rpc.constructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.security.UserCredentials; -import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.dao.util.mapping.JacksonUtil; import org.thingsboard.server.gen.edge.UpdateMsgType; +import org.thingsboard.server.gen.edge.UserCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.UserUpdateMsg; @Component @Slf4j public class UserUpdateMsgConstructor { - @Autowired - private UserService userService; - public UserUpdateMsg constructUserUpdatedMsg(UpdateMsgType msgType, User user) { UserUpdateMsg.Builder builder = UserUpdateMsg.newBuilder() .setMsgType(msgType) .setIdMSB(user.getId().getId().getMostSignificantBits()) .setIdLSB(user.getId().getId().getLeastSignificantBits()) .setEmail(user.getEmail()) - .setAuthority(user.getAuthority().name()) - .setEnabled(false); + .setAuthority(user.getAuthority().name()); if (user.getFirstName() != null) { builder.setFirstName(user.getFirstName()); } @@ -53,13 +48,6 @@ public class UserUpdateMsgConstructor { if (user.getAdditionalInfo() != null) { builder.setAdditionalInfo(JacksonUtil.toString(user.getAdditionalInfo())); } - if (msgType.equals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE) || - msgType.equals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE)) { - UserCredentials userCredentials = userService.findUserCredentialsByUserId(user.getTenantId(), user.getId()); - if (userCredentials != null) { - builder.setEnabled(userCredentials.isEnabled()).setPassword(userCredentials.getPassword()); - } - } return builder.build(); } @@ -69,4 +57,13 @@ public class UserUpdateMsgConstructor { .setIdMSB(userId.getId().getMostSignificantBits()) .setIdLSB(userId.getId().getLeastSignificantBits()).build(); } + + public UserCredentialsUpdateMsg constructUserCredentialsUpdatedMsg(UserCredentials userCredentials) { + UserCredentialsUpdateMsg.Builder builder = UserCredentialsUpdateMsg.newBuilder() + .setUserIdMSB(userCredentials.getUserId().getId().getMostSignificantBits()) + .setUserIdLSB(userCredentials.getUserId().getId().getLeastSignificantBits()) + .setEnabled(userCredentials.isEnabled()) + .setPassword(userCredentials.getPassword()); + return builder.build(); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java index 306d1bf5f9..c300a3204d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java @@ -15,10 +15,11 @@ */ package org.thingsboard.server.service.edge.rpc.init; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; import io.grpc.stub.StreamObserver; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; @@ -26,13 +27,21 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.DashboardInfo; +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.User; import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TimePageData; @@ -43,44 +52,60 @@ 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.common.data.security.DeviceCredentials; +import org.thingsboard.server.common.data.security.UserCredentials; +import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.dashboard.DashboardService; +import org.thingsboard.server.dao.device.DeviceCredentialsService; 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.dao.user.UserService; import org.thingsboard.server.gen.edge.AssetUpdateMsg; +import org.thingsboard.server.gen.edge.AttributesRequestMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.DeviceUpdateMsg; +import org.thingsboard.server.gen.edge.DownlinkMsg; +import org.thingsboard.server.gen.edge.EntityDataProto; import org.thingsboard.server.gen.edge.EntityUpdateMsg; import org.thingsboard.server.gen.edge.EntityViewUpdateMsg; +import org.thingsboard.server.gen.edge.RelationRequestMsg; 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; import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.UpdateMsgType; +import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.UserUpdateMsg; -import org.thingsboard.server.service.edge.EdgeContextComponent; import org.thingsboard.server.service.edge.rpc.constructor.AssetUpdateMsgConstructor; 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.EntityDataMsgConstructor; 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 org.thingsboard.server.service.edge.rpc.constructor.UserUpdateMsgConstructor; +import org.thingsboard.server.service.executors.DbCallbackExecutorService; import java.util.ArrayList; -import java.util.HashSet; +import java.util.Collections; import java.util.List; -import java.util.Set; import java.util.UUID; @Service @Slf4j public class DefaultSyncEdgeService implements SyncEdgeService { + private static final ObjectMapper mapper = new ObjectMapper(); + + @Autowired + private AttributesService attributesService; + @Autowired private RuleChainService ruleChainService; @@ -90,6 +115,9 @@ public class DefaultSyncEdgeService implements SyncEdgeService { @Autowired private DeviceService deviceService; + @Autowired + private DeviceCredentialsService deviceCredentialsService; + @Autowired private AssetService assetService; @@ -123,34 +151,29 @@ public class DefaultSyncEdgeService implements SyncEdgeService { @Autowired private RelationUpdateMsgConstructor relationUpdateMsgConstructor; - @Override - public void sync(EdgeContextComponent ctx, Edge edge, StreamObserver outputStream) { - Set pushedEntityIds = new HashSet<>(); - syncUsers(ctx, edge, pushedEntityIds, outputStream); - List> futures = new ArrayList<>(); - futures.add(syncRuleChains(ctx, edge, pushedEntityIds, outputStream)); - futures.add(syncDevices(ctx, edge, pushedEntityIds, outputStream)); - futures.add(syncAssets(ctx, edge, pushedEntityIds, outputStream)); - futures.add(syncEntityViews(ctx, edge, pushedEntityIds, outputStream)); - futures.add(syncDashboards(ctx, edge, pushedEntityIds, outputStream)); - Futures.addCallback(Futures.allAsList(futures), new FutureCallback>() { - @Override - public void onSuccess(@Nullable List result) { - syncRelations(ctx, edge, pushedEntityIds, outputStream); - } + @Autowired + private EntityDataMsgConstructor entityDataMsgConstructor; - @Override - public void onFailure(Throwable t) { - log.warn("Exception during sync entities", t); - } - }, MoreExecutors.directExecutor()); + @Autowired + private DbCallbackExecutorService dbCallbackExecutorService; + + @Override + public void sync(Edge edge, StreamObserver outputStream) { + syncUsers(edge, outputStream); + syncRuleChains(edge, outputStream); + syncDevices(edge, outputStream); + syncAssets(edge, outputStream); + syncEntityViews(edge, outputStream); + syncDashboards(edge, outputStream); } - private ListenableFuture syncRuleChains(EdgeContextComponent ctx, Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + private void syncRuleChains(Edge edge, StreamObserver outputStream) { try { - ListenableFuture> future = ruleChainService.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); - return Futures.transform(future, pageData -> { - try { + ListenableFuture> future = + ruleChainService.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); + Futures.addCallback(future, new FutureCallback>() { + @Override + public void onSuccess(@Nullable TimePageData pageData) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { log.trace("[{}] [{}] rule chains(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); for (RuleChain ruleChain : pageData.getData()) { @@ -165,81 +188,93 @@ public class DefaultSyncEdgeService implements SyncEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); - pushedEntityIds.add(ruleChain.getId()); } } - } catch (Exception e) { - log.error("Exception during loading edge rule chain(s) on sync!", e); } - return null; - }, ctx.getDbCallbackExecutor()); + + @Override + public void onFailure(Throwable t) { + log.error("Exception during loading edge rule chain(s) on sync!", t); + } + }, dbCallbackExecutorService); } catch (Exception e) { log.error("Exception during loading edge rule chain(s) on sync!", e); - return Futures.immediateFuture(null); } } - private ListenableFuture syncDevices(EdgeContextComponent ctx, Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + private void syncDevices(Edge edge, StreamObserver outputStream) { try { - ListenableFuture> future = deviceService.findDevicesByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); - return Futures.transform(future, pageData -> { - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - log.trace("[{}] [{}] device(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); - for (Device device : pageData.getData()) { - DeviceUpdateMsg deviceUpdateMsg = - deviceUpdateMsgConstructor.constructDeviceUpdatedMsg( - UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, - device); - EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() - .setDeviceUpdateMsg(deviceUpdateMsg) - .build(); - outputStream.onNext(ResponseMsg.newBuilder() - .setEntityUpdateMsg(entityUpdateMsg) - .build()); - pushedEntityIds.add(device.getId()); + ListenableFuture> future = + deviceService.findDevicesByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); + Futures.addCallback(future, new FutureCallback>() { + @Override + public void onSuccess(@Nullable TimePageData pageData) { + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + log.trace("[{}] [{}] device(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); + for (Device device : pageData.getData()) { + DeviceUpdateMsg deviceUpdateMsg = + deviceUpdateMsgConstructor.constructDeviceUpdatedMsg( + UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, + device); + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setDeviceUpdateMsg(deviceUpdateMsg) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } } } - return null; - }, ctx.getDbCallbackExecutor()); + + @Override + public void onFailure(Throwable t) { + log.error("Exception during loading edge device(s) on sync!", t); + } + }, dbCallbackExecutorService); } catch (Exception e) { log.error("Exception during loading edge device(s) on sync!", e); - return Futures.immediateFuture(null); } } - private ListenableFuture syncAssets(EdgeContextComponent ctx, Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + private void syncAssets(Edge edge, StreamObserver outputStream) { try { ListenableFuture> future = assetService.findAssetsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); - return Futures.transform(future, pageData -> { - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - log.trace("[{}] [{}] asset(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); - for (Asset asset : pageData.getData()) { - AssetUpdateMsg assetUpdateMsg = - assetUpdateMsgConstructor.constructAssetUpdatedMsg( - UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, - asset); - EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() - .setAssetUpdateMsg(assetUpdateMsg) - .build(); - outputStream.onNext(ResponseMsg.newBuilder() - .setEntityUpdateMsg(entityUpdateMsg) - .build()); - pushedEntityIds.add(asset.getId()); + Futures.addCallback(future, new FutureCallback>() { + @Override + public void onSuccess(@Nullable TimePageData pageData) { + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + log.trace("[{}] [{}] asset(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); + for (Asset asset : pageData.getData()) { + AssetUpdateMsg assetUpdateMsg = + assetUpdateMsgConstructor.constructAssetUpdatedMsg( + UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, + asset); + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setAssetUpdateMsg(assetUpdateMsg) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } } } - return null; - }, ctx.getDbCallbackExecutor()); + + @Override + public void onFailure(Throwable t) { + log.error("Exception during loading edge asset(s) on sync!", t); + } + }, dbCallbackExecutorService); } catch (Exception e) { log.error("Exception during loading edge asset(s) on sync!", e); - return Futures.immediateFuture(null); } } - private ListenableFuture syncEntityViews(EdgeContextComponent ctx, Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + private void syncEntityViews(Edge edge, StreamObserver outputStream) { try { ListenableFuture> future = entityViewService.findEntityViewsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); - return Futures.transform(future, pageData -> { - try { + Futures.addCallback(future, new FutureCallback>() { + @Override + public void onSuccess(@Nullable TimePageData pageData) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { log.trace("[{}] [{}] entity view(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); for (EntityView entityView : pageData.getData()) { @@ -253,25 +288,26 @@ public class DefaultSyncEdgeService implements SyncEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); - pushedEntityIds.add(entityView.getId()); } } - } catch (Exception e) { - log.error("Exception during loading edge entity view(s) on sync!", e); } - return null; - }, ctx.getDbCallbackExecutor()); + + @Override + public void onFailure(Throwable t) { + log.error("Exception during loading edge entity view(s) on sync!", t); + } + }, dbCallbackExecutorService); } catch (Exception e) { log.error("Exception during loading edge entity view(s) on sync!", e); - return Futures.immediateFuture(null); } } - private ListenableFuture syncDashboards(EdgeContextComponent ctx, Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + private void syncDashboards(Edge edge, StreamObserver outputStream) { try { ListenableFuture> future = dashboardService.findDashboardsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); - return Futures.transform(future, pageData -> { - try { + Futures.addCallback(future, new FutureCallback>() { + @Override + public void onSuccess(@Nullable TimePageData pageData) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { log.trace("[{}] [{}] dashboard(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); for (DashboardInfo dashboardInfo : pageData.getData()) { @@ -286,34 +322,34 @@ public class DefaultSyncEdgeService implements SyncEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); - pushedEntityIds.add(dashboard.getId()); } } - } catch (Exception e) { - log.error("Exception during loading edge dashboard(s) on sync!", e); } - return null; - }, ctx.getDbCallbackExecutor()); + + @Override + public void onFailure(Throwable t) { + log.error("Exception during loading edge dashboard(s) on sync!", t); + } + }, dbCallbackExecutorService); } catch (Exception e) { log.error("Exception during loading edge dashboard(s) on sync!", e); - return Futures.immediateFuture(null); } } - private void syncUsers(EdgeContextComponent ctx, Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + private void syncUsers(Edge edge, StreamObserver outputStream) { try { TextPageData pageData = userService.findTenantAdmins(edge.getTenantId(), new TextPageLink(Integer.MAX_VALUE)); - pushUsersToEdge(pageData, edge, pushedEntityIds, outputStream); + pushUsersToEdge(pageData, edge, outputStream); if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { pageData = userService.findCustomerUsers(edge.getTenantId(), edge.getCustomerId(), new TextPageLink(Integer.MAX_VALUE)); - pushUsersToEdge(pageData, edge, pushedEntityIds, outputStream); + pushUsersToEdge(pageData, edge, outputStream); } } catch (Exception e) { log.error("Exception during loading edge user(s) on sync!", e); } } - private void pushUsersToEdge(TextPageData pageData, Edge edge, Set pushedEntityIds, StreamObserver outputStream) { + private void pushUsersToEdge(TextPageData pageData, Edge edge, StreamObserver outputStream) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { log.trace("[{}] [{}] user(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); for (User user : pageData.getData()) { @@ -327,81 +363,161 @@ public class DefaultSyncEdgeService implements SyncEdgeService { outputStream.onNext(ResponseMsg.newBuilder() .setEntityUpdateMsg(entityUpdateMsg) .build()); - pushedEntityIds.add(user.getId()); } } } - private ListenableFuture syncRelations(EdgeContextComponent ctx, 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)); + @Override + public void processRuleChainMetadata(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); + RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = + ruleChainUpdateMsgConstructor.constructRuleChainMetadataUpdatedMsg( + UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, + ruleChainMetaData); + if (ruleChainMetadataUpdateMsg != null) { + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setRuleChainMetadataUpdateMsg(ruleChainMetadataUpdateMsg) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } + } + } + + @Override + public void processAttributesRequestMsg(Edge edge, AttributesRequestMsg attributesRequestMsg, StreamObserver outputStream) { + EntityId entityId = EntityIdFactory.getByTypeAndUuid( + EntityType.valueOf(attributesRequestMsg.getEntityType()), + new UUID(attributesRequestMsg.getEntityIdMSB(), attributesRequestMsg.getEntityIdLSB())); + ListenableFuture> ssAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SERVER_SCOPE); + Futures.addCallback(ssAttrFuture, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List ssAttributes) { + if (ssAttributes != null && !ssAttributes.isEmpty()) { + try { + ObjectNode entityNode = mapper.createObjectNode(); + for (AttributeKvEntry attr : ssAttributes) { + if (attr.getDataType() == DataType.BOOLEAN && attr.getBooleanValue().isPresent()) { + entityNode.put(attr.getKey(), attr.getBooleanValue().get()); + } else if (attr.getDataType() == DataType.DOUBLE && attr.getDoubleValue().isPresent()) { + entityNode.put(attr.getKey(), attr.getDoubleValue().get()); + } else if (attr.getDataType() == DataType.LONG && attr.getLongValue().isPresent()) { + entityNode.put(attr.getKey(), attr.getLongValue().get()); + } else { + entityNode.put(attr.getKey(), attr.getValueAsString()); + } + } + log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, entityNode); + + EntityDataProto entityDataProto = + entityDataMsgConstructor.constructEntityDataMsg( + entityId, + ActionType.ATTRIBUTES_UPDATED, + JsonUtils.parse(mapper.writeValueAsString(entityNode))); + DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() + .addAllEntityData(Collections.singletonList(entityDataProto)); + DownlinkMsg value = builder.build(); + + outputStream.onNext(ResponseMsg.newBuilder() + .setDownlinkMsg(value).build()); + } catch (Exception e) { + log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e); + } + } + } + + @Override + public void onFailure(Throwable t) { + } - ListenableFuture>> relationsListFuture = Futures.allAsList(futures); - return Futures.transform(relationsListFuture, relationsList -> { + }, dbCallbackExecutorService); + + // TODO: voba - push shared attributes to edge? + ListenableFuture> shAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SHARED_SCOPE); + ListenableFuture> clAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.CLIENT_SCOPE); + } + + @Override + public void processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg, StreamObserver outputStream) { + EntityId entityId = EntityIdFactory.getByTypeAndUuid( + EntityType.valueOf(relationRequestMsg.getEntityType()), + new UUID(relationRequestMsg.getEntityIdMSB(), relationRequestMsg.getEntityIdLSB())); + + List>> futures = new ArrayList<>(); + futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.FROM)); + futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.TO)); + ListenableFuture>> relationsListFuture = Futures.allAsList(futures); + Futures.addCallback(relationsListFuture, new FutureCallback>>() { + @Override + public void onSuccess(@Nullable List> 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 sync!", relation, e); + log.trace("[{}] [{}] [{}] relation(s) are going to be pushed to edge.", edge.getId(), entityId, entityRelations.size()); + for (EntityRelation relation : entityRelations) { + 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 sync!", relation, e); + } } } } } catch (Exception e) { log.error("Exception during loading relation(s) to edge on sync!", e); } - return null; - }, ctx.getDbCallbackExecutor()); - } else { - return Futures.immediateFuture(null); - } + } + + @Override + public void onFailure(Throwable t) { + log.error("Exception during loading relation(s) to edge on sync!", t); + } + }, dbCallbackExecutorService); } - private ListenableFuture> syncRelations(Edge edge, EntityId entityId, EntitySearchDirection direction) { + private ListenableFuture> findRelationByQuery(Edge edge, EntityId entityId, EntitySearchDirection direction) { EntityRelationsQuery query = new EntityRelationsQuery(); query.setParameters(new RelationsSearchParameters(entityId, direction, -1, false)); return relationService.findByQuery(edge.getTenantId(), query); } @Override - 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); - RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = - ruleChainUpdateMsgConstructor.constructRuleChainMetadataUpdatedMsg( - UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, - ruleChainMetaData); - if (ruleChainMetadataUpdateMsg != null) { - EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() - .setRuleChainMetadataUpdateMsg(ruleChainMetadataUpdateMsg) - .build(); - outputStream.onNext(ResponseMsg.newBuilder() - .setEntityUpdateMsg(entityUpdateMsg) - .build()); - } + public void processDeviceCredentialsRequestMsg(Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg, StreamObserver outputStream) { + DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsRequestMsg.getDeviceIdMSB(), deviceCredentialsRequestMsg.getDeviceIdLSB())); + DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(edge.getTenantId(), deviceId); + if (deviceCredentials != null) { + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setDeviceCredentialsUpdateMsg(deviceUpdateMsgConstructor.constructDeviceCredentialsUpdatedMsg(deviceCredentials)) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } + } + + @Override + public void processUserCredentialsRequestMsg(Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg, StreamObserver outputStream) { + UserId userId = new UserId(new UUID(userCredentialsRequestMsg.getUserIdMSB(), userCredentialsRequestMsg.getUserIdLSB())); + UserCredentials userCredentialsByUserId = userService.findUserCredentialsByUserId(edge.getTenantId(), userId); + if (userCredentialsByUserId != null) { + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setUserCredentialsUpdateMsg(userUpdateMsgConstructor.constructUserCredentialsUpdatedMsg(userCredentialsByUserId)) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); } } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java index 84b4d9ffff..dbcb659daf 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java @@ -17,13 +17,24 @@ package org.thingsboard.server.service.edge.rpc.init; import io.grpc.stub.StreamObserver; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.gen.edge.AttributesRequestMsg; +import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg; +import org.thingsboard.server.gen.edge.RelationRequestMsg; import org.thingsboard.server.gen.edge.ResponseMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg; -import org.thingsboard.server.service.edge.EdgeContextComponent; +import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg; public interface SyncEdgeService { - void sync(EdgeContextComponent ctx, Edge edge, StreamObserver outputStream); + void sync(Edge edge, StreamObserver outputStream); - void syncRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver outputStream); + void processRuleChainMetadata(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg, StreamObserver outputStream); + + void processAttributesRequestMsg(Edge edge, AttributesRequestMsg attributesRequestMsg, StreamObserver outputStream); + + void processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg, StreamObserver outputStream); + + void processDeviceCredentialsRequestMsg(Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg, StreamObserver outputStream); + + void processUserCredentialsRequestMsg(Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg, StreamObserver outputStream); } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 38cbe7092c..8c33a55a6e 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -48,15 +48,17 @@ message ResponseMsg { message EntityUpdateMsg { DeviceUpdateMsg deviceUpdateMsg = 1; - RuleChainUpdateMsg ruleChainUpdateMsg = 2; - RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = 3; - DashboardUpdateMsg dashboardUpdateMsg = 4; - AssetUpdateMsg assetUpdateMsg = 5; - EntityViewUpdateMsg entityViewUpdateMsg = 6; - AlarmUpdateMsg alarmUpdateMsg = 7; - UserUpdateMsg userUpdateMsg = 8; - CustomerUpdateMsg customerUpdateMsg = 9; - RelationUpdateMsg relationUpdateMsg = 10; + DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = 2; + RuleChainUpdateMsg ruleChainUpdateMsg = 3; + RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = 4; + DashboardUpdateMsg dashboardUpdateMsg = 5; + AssetUpdateMsg assetUpdateMsg = 6; + EntityViewUpdateMsg entityViewUpdateMsg = 7; + AlarmUpdateMsg alarmUpdateMsg = 8; + UserUpdateMsg userUpdateMsg = 9; + UserCredentialsUpdateMsg userCredentialsUpdateMsg = 10; + CustomerUpdateMsg customerUpdateMsg = 11; + RelationUpdateMsg relationUpdateMsg = 12; } enum RequestMsgType { @@ -169,9 +171,14 @@ message DeviceUpdateMsg { string name = 4; string type = 5; string label = 6; - string credentialsType = 7; - string credentialsId = 8; - string credentialsValue = 9; +} + +message DeviceCredentialsUpdateMsg { + int64 deviceIdMSB = 1; + int64 deviceIdLSB = 2; + string credentialsType = 3; + string credentialsId = 4; + string credentialsValue = 5; } message AssetUpdateMsg { @@ -248,8 +255,13 @@ message UserUpdateMsg { string firstName = 6; string lastName = 7; string additionalInfo = 8; - bool enabled = 9; - string password = 10; +} + +message UserCredentialsUpdateMsg { + int64 userIdMSB = 1; + int64 userIdLSB = 2; + bool enabled = 3; + string password = 4; } message RuleChainMetadataRequestMsg { @@ -257,6 +269,28 @@ message RuleChainMetadataRequestMsg { int64 ruleChainIdLSB = 2; } +message AttributesRequestMsg { + int64 entityIdMSB = 1; + int64 entityIdLSB = 2; + string entityType = 3; +} + +message RelationRequestMsg { + int64 entityIdMSB = 1; + int64 entityIdLSB = 2; + string entityType = 3; +} + +message UserCredentialsRequestMsg { + int64 userIdMSB = 1; + int64 userIdLSB = 2; +} + +message DeviceCredentialsRequestMsg { + int64 deviceIdMSB = 1; + int64 deviceIdLSB = 2; +} + enum EdgeEntityType { DEVICE = 0; ASSET = 1; @@ -270,8 +304,13 @@ message UplinkMsg { int32 uplinkMsgId = 1; repeated EntityDataProto entityData = 2; repeated DeviceUpdateMsg deviceUpdateMsg = 3; - repeated AlarmUpdateMsg alarmUpdateMsg = 4; - repeated RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg = 5; + repeated DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = 4; + repeated AlarmUpdateMsg alarmUpdateMsg = 5; + repeated RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg = 6; + repeated AttributesRequestMsg attributesRequestMsg = 7; + repeated RelationRequestMsg relationRequestMsg = 8; + repeated UserCredentialsRequestMsg userCredentialsRequestMsg = 9; + repeated DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = 10; } message UplinkResponseMsg { @@ -282,5 +321,6 @@ message UplinkResponseMsg { message DownlinkMsg { int32 downlinkMsgId = 1; repeated EntityDataProto entityData = 2; + repeated DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = 3; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java index ae3f787e5a..06bf99f5c3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java @@ -232,7 +232,7 @@ public class DashboardServiceImpl extends AbstractEntityService implements Dashb if (edge == null) { throw new DataValidationException("Can't assign dashboard to non-existent edge!"); } - if (!edge.getTenantId().getId().equals(dashboard.getTenantId().getId())) { + if (!edge.getTenantId().equals(dashboard.getTenantId())) { throw new DataValidationException("Can't assign dashboard to edge from different tenant!"); } try { diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java index fbe44c2267..ec99c673a6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java @@ -22,17 +22,13 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntitySubtype; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.DashboardId; -import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageLink; -import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; -import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.dao.model.nosql.EdgeEntity; import org.thingsboard.server.dao.nosql.CassandraAbstractSearchTextDao; import org.thingsboard.server.dao.relation.RelationDao; @@ -113,19 +109,17 @@ public class CassandraEdgeDao extends CassandraAbstractSearchTextDao> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId) { log.debug("Try to find edges by tenantId [{}], ruleChainId [{}]", tenantId, ruleChainId); ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return Futures.transformAsync(relations, input -> { - List> edgeFutures = new ArrayList<>(input.size()); - for (EntityRelation relation : input) { - edgeFutures.add(findByIdAsync(new TenantId(tenantId), relation.getFrom().getId())); - } - return Futures.successfulAsList(edgeFutures); - }, MoreExecutors.directExecutor()); + return transformFromRelationToEdge(tenantId, relations); } @Override public ListenableFuture> findEdgesByTenantIdAndDashboardId(UUID tenantId, UUID dashboardId) { log.debug("Try to find edges by tenantId [{}], dashboardId [{}]", tenantId, dashboardId); ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new DashboardId(dashboardId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return transformFromRelationToEdge(tenantId, relations); + } + + private ListenableFuture> transformFromRelationToEdge(UUID tenantId, ListenableFuture> relations) { return Futures.transformAsync(relations, input -> { List> edgeFutures = new ArrayList<>(input.size()); for (EntityRelation relation : input) { 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 d07f0704ca..73272c3328 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 @@ -414,7 +414,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC if (edge == null) { throw new DataValidationException("Can't assign ruleChain to non-existent edge!"); } - if (!edge.getTenantId().getId().equals(ruleChain.getTenantId().getId())) { + if (!edge.getTenantId().equals(ruleChain.getTenantId())) { throw new DataValidationException("Can't assign ruleChain to edge from different tenant!"); } try { 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 6a18a4e039..feacf23f5e 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 @@ -151,19 +151,17 @@ public class JpaEdgeDao extends JpaAbstractSearchTextDao imple public ListenableFuture> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId) { log.debug("Try to find edges by tenantId [{}], ruleChainId [{}]", tenantId, ruleChainId); ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return Futures.transformAsync(relations, input -> { - List> edgeFutures = new ArrayList<>(input.size()); - for (EntityRelation relation : input) { - edgeFutures.add(findByIdAsync(new TenantId(tenantId), relation.getFrom().getId())); - } - return Futures.successfulAsList(edgeFutures); - }, MoreExecutors.directExecutor()); + return transformFromRelationToEdge(tenantId, relations); } @Override public ListenableFuture> findEdgesByTenantIdAndDashboardId(UUID tenantId, UUID dashboardId) { log.debug("Try to find edges by tenantId [{}], dashboardId [{}]", tenantId, dashboardId); ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new DashboardId(dashboardId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return transformFromRelationToEdge(tenantId, relations); + } + + private ListenableFuture> transformFromRelationToEdge(UUID tenantId, ListenableFuture> relations) { return Futures.transformAsync(relations, input -> { List> edgeFutures = new ArrayList<>(input.size()); for (EntityRelation relation : input) {