diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index 279ec96064..27489c98cc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -242,7 +242,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { } } } else { - ListenableFuture> edgeIdsFuture = findRelatedEdgeIdsByEntityId(tenantId, entityId); + ListenableFuture> edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, dbCallbackExecutorService); Futures.transform(edgeIdsFuture, edgeIds -> { if (edgeIds != null && !edgeIds.isEmpty()) { for (EdgeId edgeId : edgeIds) { @@ -321,7 +321,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { if (alarm != null) { EdgeEventType edgeEventType = getEdgeQueueTypeByEntityType(alarm.getOriginator().getEntityType()); if (edgeEventType != null) { - ListenableFuture> relatedEdgeIdsByEntityIdFuture = findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator()); + ListenableFuture> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), dbCallbackExecutorService); Futures.transform(relatedEdgeIdsByEntityIdFuture, relatedEdgeIdsByEntityId -> { if (relatedEdgeIdsByEntityId != null) { for (EdgeId edgeId : relatedEdgeIdsByEntityId) { @@ -346,8 +346,8 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && !relation.getTo().getEntityType().equals(EntityType.EDGE)) { List>> futures = new ArrayList<>(); - futures.add(findRelatedEdgeIdsByEntityId(tenantId, relation.getTo())); - futures.add(findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom())); + futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo(), dbCallbackExecutorService)); + futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom(), dbCallbackExecutorService)); ListenableFuture>> combinedFuture = Futures.allAsList(futures); Futures.transform(combinedFuture, listOfListsEdgeIds -> { Set uniqueEdgeIds = new HashSet<>(); @@ -373,48 +373,6 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { } } - private ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) { - switch (entityId.getEntityType()) { - case DEVICE: - case ASSET: - case ENTITY_VIEW: - ListenableFuture> originatorEdgeRelationsFuture = - relationService.findByToAndTypeAsync(tenantId, entityId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return Futures.transform(originatorEdgeRelationsFuture, originatorEdgeRelations -> { - if (originatorEdgeRelations != null && originatorEdgeRelations.size() > 0) { - return Collections.singletonList(new EdgeId(originatorEdgeRelations.get(0).getFrom().getId())); - } else { - return Collections.emptyList(); - } - }, dbCallbackExecutorService); - case DASHBOARD: - return convertToEdgeIds(edgeService.findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId()))); - case RULE_CHAIN: - return convertToEdgeIds(edgeService.findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId()))); - case USER: - User userById = userService.findUserById(tenantId, new UserId(entityId.getId())); - TextPageData edges; - if (userById.getCustomerId() == null || userById.getCustomerId().getId().equals(ModelConstants.NULL_UUID)) { - edges = edgeService.findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE)); - } else { - edges = edgeService.findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), new TextPageLink(Integer.MAX_VALUE)); - } - return convertToEdgeIds(Futures.immediateFuture(edges.getData())); - default: - return Futures.immediateFuture(Collections.emptyList()); - } - } - - private ListenableFuture> convertToEdgeIds(ListenableFuture> future) { - return Futures.transform(future, edges -> { - if (edges != null && !edges.isEmpty()) { - return edges.stream().map(IdBased::getId).collect(Collectors.toList()); - } else { - return Collections.emptyList(); - } - }, dbCallbackExecutorService); - } - private EdgeEventType getEdgeQueueTypeByEntityType(EntityType entityType) { switch (entityType) { case DEVICE: diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 121772a01b..86c45c78ba 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -33,6 +33,7 @@ import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang.RandomStringUtils; import org.checkerframework.checker.nullness.qual.Nullable; +import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; @@ -64,6 +65,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetsBundleId; +import org.thingsboard.server.common.data.kv.AttributeKey; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; @@ -88,6 +90,7 @@ import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.gen.edge.AdminSettingsUpdateMsg; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg; +import org.thingsboard.server.gen.edge.AttributeDeleteMsg; import org.thingsboard.server.gen.edge.AttributesRequestMsg; import org.thingsboard.server.gen.edge.ConnectRequestMsg; import org.thingsboard.server.gen.edge.ConnectResponseCode; @@ -125,11 +128,12 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; import java.io.Closeable; -import java.io.IOException; import java.util.ArrayList; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Optional; +import java.util.Set; import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; @@ -382,7 +386,7 @@ public final class EdgeGrpcSession implements Closeable { ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); } - private DownlinkMsg processTelemetryMessage(EdgeEvent edgeEvent) throws IOException { + private DownlinkMsg processTelemetryMessage(EdgeEvent edgeEvent) { log.trace("Executing processTelemetryMessage, edgeEvent [{}]", edgeEvent); EntityId entityId = null; switch (edgeEvent.getEdgeEventType()) { @@ -823,6 +827,9 @@ public final class EdgeGrpcSession implements Closeable { result.add(processPostTelemetry(entityId, entityData.getPostTelemetryMsg(), metaData)); } } + if (entityData.hasAttributeDeleteMsg()) { + result.add(processAttributeDeleteMsg(entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType())); + } } } @@ -948,6 +955,26 @@ public final class EdgeGrpcSession implements Closeable { return Futures.immediateFuture(null); } + private ListenableFuture processAttributeDeleteMsg(EntityId entityId, AttributeDeleteMsg attributeDeleteMsg, String entityType) { + try { + String scope = attributeDeleteMsg.getScope(); + List attributeNames = attributeDeleteMsg.getAttributeNamesList(); + ctx.getAttributesService().removeAll(edge.getTenantId(), entityId, scope, attributeNames); + if (EntityType.DEVICE.name().equals(entityType)) { + Set attributeKeys = new HashSet<>(); + for (String attributeName : attributeNames) { + attributeKeys.add(new AttributeKey(scope, attributeName)); + } + ctx.getTbClusterService().pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( + edge.getTenantId(), (DeviceId) entityId, attributeKeys), null); + } + } catch (Exception e) { + log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, e); + return Futures.immediateFailedFuture(new RuntimeException("Can't process attribute delete msg " + attributeDeleteMsg, e)); + } + return Futures.immediateFuture(null); + } + private ListenableFuture onDeviceUpdate(DeviceUpdateMsg deviceUpdateMsg) { DeviceId edgeDeviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); switch (deviceUpdateMsg.getMsgType()) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java index c1d211af58..0f7fb0523d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java @@ -43,9 +43,11 @@ public class EntityDataMsgConstructor { case TIMESERIES_UPDATED: try { JsonObject data = entityData.getAsJsonObject(); - long ts = System.currentTimeMillis(); + long ts; if (data.get("ts") != null && !data.get("ts").isJsonNull()) { - ts = data.getAsJsonObject("ts").getAsLong(); + ts = data.getAsJsonPrimitive("ts").getAsLong(); + } else { + ts = System.currentTimeMillis(); } builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data.getAsJsonObject("data"), ts)); } catch (Exception e) { diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index 5c6e5d7559..71c56068d1 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -22,15 +22,15 @@ import org.thingsboard.server.common.data.edge.EdgeSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; -import org.thingsboard.server.common.data.page.TimePageData; -import org.thingsboard.server.common.data.page.TimePageLink; import java.util.List; import java.util.Optional; +import java.util.concurrent.Executor; public interface EdgeService { @@ -75,6 +75,8 @@ public interface EdgeService { ListenableFuture> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId); ListenableFuture> findEdgesByTenantIdAndDashboardId(TenantId tenantId, DashboardId dashboardId); + + ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, Executor executorService); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index 01f82faa17..4ce86beea5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -31,31 +31,35 @@ import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.IdBased; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; -import org.thingsboard.server.common.data.page.TimePageData; -import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; +import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.dao.customer.CustomerDao; import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.Validator; import org.thingsboard.server.dao.tenant.TenantDao; +import org.thingsboard.server.dao.user.UserService; import javax.annotation.Nullable; import java.util.ArrayList; @@ -63,6 +67,7 @@ import java.util.Collections; import java.util.Comparator; import java.util.List; import java.util.Optional; +import java.util.concurrent.Executor; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.CacheConstants.EDGE_CACHE; @@ -91,6 +96,9 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic @Autowired private CustomerDao customerDao; + @Autowired + private UserService userService; + @Autowired private CacheManager cacheManager; @@ -420,4 +428,47 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic } }; + @Override + public ListenableFuture> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, Executor executorService) { + switch (entityId.getEntityType()) { + case DEVICE: + case ASSET: + case ENTITY_VIEW: + ListenableFuture> originatorEdgeRelationsFuture = + relationService.findByToAndTypeAsync(tenantId, entityId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return Futures.transform(originatorEdgeRelationsFuture, originatorEdgeRelations -> { + if (originatorEdgeRelations != null && originatorEdgeRelations.size() > 0) { + return Collections.singletonList(new EdgeId(originatorEdgeRelations.get(0).getFrom().getId())); + } else { + return Collections.emptyList(); + } + }, executorService); + case DASHBOARD: + return convertToEdgeIds(findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId())), executorService); + case RULE_CHAIN: + return convertToEdgeIds(findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId())), executorService); + case USER: + User userById = userService.findUserById(tenantId, new UserId(entityId.getId())); + TextPageData edges; + if (userById.getCustomerId() == null || userById.getCustomerId().getId().equals(ModelConstants.NULL_UUID)) { + edges = findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE)); + } else { + edges = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), new TextPageLink(Integer.MAX_VALUE)); + } + return convertToEdgeIds(Futures.immediateFuture(edges.getData()), executorService); + default: + return Futures.immediateFuture(Collections.emptyList()); + } + } + + private ListenableFuture> convertToEdgeIds(ListenableFuture> future, Executor executorService) { + return Futures.transform(future, edges -> { + if (edges != null && !edges.isEmpty()) { + return edges.stream().map(IdBased::getId).collect(Collectors.toList()); + } else { + return Collections.emptyList(); + } + }, executorService); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java index be265ebc19..b531af7ed3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java @@ -232,14 +232,7 @@ public class TbMsgPushToEdgeNode implements TbNode { TextPageData edgesByTenantId = ctx.getEdgeService().findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE)); return Futures.immediateFuture(edgesByTenantId.getData().stream().map(IdBased::getId).collect(Collectors.toList())); } else { - ListenableFuture> future = ctx.getRelationService().findByToAndTypeAsync(tenantId, originatorId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); - return Futures.transform(future, relations -> { - List result = new ArrayList<>(); - if (relations != null && relations.size() > 0) { - result.add(new EdgeId(relations.get(0).getFrom().getId())); - } - return result; - }, ctx.getDbCallbackExecutor()); + return ctx.getEdgeService().findRelatedEdgeIdsByEntityId(tenantId, originatorId, ctx.getDbCallbackExecutor()); } }