Browse Source

Pagedata usage for get edges methods

pull/4918/head
Volodymyr Babak 5 years ago
parent
commit
b0bfdfff8e
  1. 22
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  2. 4
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  3. 40
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java
  4. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java
  6. 60
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java
  7. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java
  8. 83
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java
  9. 62
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java
  10. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java
  11. 6
      common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java
  12. 18
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java
  13. 117
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  14. 4
      dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java
  15. 22
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  16. 11
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java
  17. 23
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java
  18. 105
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

22
application/src/main/java/org/thingsboard/server/controller/BaseController.java

@ -83,6 +83,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.DataType;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.page.SortOrder;
import org.thingsboard.server.common.data.page.TimePageLink;
@ -145,6 +146,7 @@ import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import javax.mail.MessagingException;
import javax.servlet.http.HttpServletResponse;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
@ -164,6 +166,8 @@ public abstract class BaseController {
protected static final String DEFAULT_DASHBOARD = "defaultDashboardId";
protected static final String HOME_DASHBOARD = "homeDashboardId";
private static final int DEFAULT_PAGE_SIZE = 1000;
private static final ObjectMapper json = new ObjectMapper();
@Autowired
@ -1082,12 +1086,18 @@ public abstract class BaseController {
if (!edgesEnabled) {
return null;
}
List<EdgeId> result = null;
try {
result = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId).get();
} catch (Exception e) {
log.error("[{}] can't find related edge ids for entity [{}]", tenantId, entityId, e);
}
List<EdgeId> result = new ArrayList<>();
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<EdgeId> pageData;
do {
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
result.addAll(pageData.getData());
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
return result;
}

4
application/src/main/java/org/thingsboard/server/controller/RuleChainController.java

@ -86,7 +86,7 @@ public class RuleChainController extends BaseController {
public static final String RULE_CHAIN_ID = "ruleChainId";
public static final String RULE_NODE_ID = "ruleNodeId";
private static final int DEFAULT_LIMIT = 100;
private static final int DEFAULT_PAGE_SIZE = 1000;
private static final ObjectMapper objectMapper = new ObjectMapper();
@ -643,7 +643,7 @@ public class RuleChainController extends BaseController {
try {
TenantId tenantId = getCurrentUser().getTenantId();
List<RuleChain> result = new ArrayList<>();
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<RuleChain> pageData;
do {
pageData = ruleChainService.findAutoAssignToEdgeRuleChainsByTenantId(tenantId, pageLink);

40
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java

@ -34,6 +34,8 @@ import org.thingsboard.server.common.data.id.AlarmId;
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.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.DownlinkMsg;
import org.thingsboard.server.gen.edge.UpdateMsgType;
@ -41,7 +43,6 @@ import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.TbCoreComponent;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
@Component
@ -138,34 +139,31 @@ public class AlarmEdgeProcessor extends BaseEdgeProcessor {
if (alarm != null) {
EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType());
if (type != null) {
ListenableFuture<List<EdgeId>> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator());
Futures.addCallback(relatedEdgeIdsByEntityIdFuture, new FutureCallback<List<EdgeId>>() {
@Override
public void onSuccess(@Nullable List<EdgeId> relatedEdgeIdsByEntityId) {
if (relatedEdgeIdsByEntityId != null) {
for (EdgeId edgeId : relatedEdgeIdsByEntityId) {
saveEdgeEvent(tenantId,
edgeId,
EdgeEventType.ALARM,
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()),
alarmId,
null);
}
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<EdgeId> pageData;
do {
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (EdgeId edgeId : pageData.getData()) {
saveEdgeEvent(tenantId,
edgeId,
EdgeEventType.ALARM,
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()),
alarmId,
null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] can't find related edge ids by entity id [{}]", tenantId.getId(), alarm.getOriginator(), t);
}
}, dbCallbackExecutorService);
} while (pageData != null && pageData.hasNext());
}
}
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] can't find alarm by id [{}]", tenantId.getId(), alarmId.getId(), t);
log.warn("[{}] can't find alarm by id [{}] {}", tenantId.getId(), alarmId.getId(), t);
}
}, dbCallbackExecutorService);
}

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java

@ -74,7 +74,7 @@ public abstract class BaseEdgeProcessor {
protected static final ObjectMapper mapper = new ObjectMapper();
protected static final int DEFAULT_LIMIT = 100;
protected static final int DEFAULT_PAGE_SIZE = 1000;
@Autowired
protected RuleChainService ruleChainService;
@ -221,7 +221,7 @@ public abstract class BaseEdgeProcessor {
}
protected void processActionForAllEdges(TenantId tenantId, EdgeEventType type, EdgeEventActionType actionType, EntityId entityId) {
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<Edge> pageData;
do {
pageData = edgeService.findEdgesByTenantId(tenantId, pageLink);

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java

@ -75,7 +75,7 @@ public class CustomerEdgeProcessor extends BaseEdgeProcessor {
CustomerId customerId = new CustomerId(EntityIdFactory.getByEdgeEventTypeAndUuid(type, uuid).getId());
switch (actionType) {
case UPDATED:
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<Edge> pageData;
do {
pageData = edgeService.findEdgesByTenantIdAndCustomerId(tenantId, customerId, pageLink);

60
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java

@ -17,14 +17,12 @@ package org.thingsboard.server.service.edge.rpc.processor;
import com.fasterxml.jackson.core.JsonProcessingException;
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.SettableFuture;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.RandomStringUtils;
import org.apache.commons.lang3.StringUtils;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.stereotype.Component;
import org.thingsboard.rule.engine.api.RpcError;
import org.thingsboard.server.common.data.Customer;
@ -40,6 +38,8 @@ import org.thingsboard.server.common.data.id.DeviceProfileId;
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.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.security.DeviceCredentials;
@ -61,7 +61,6 @@ import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.locks.ReentrantLock;
@ -79,41 +78,34 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
String deviceName = deviceUpdateMsg.getName();
Device device = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName);
if (device != null) {
ListenableFuture<List<EdgeId>> future = edgeService.findRelatedEdgeIdsByEntityId(tenantId, device.getId());
SettableFuture<Void> futureToSet = SettableFuture.create();
Futures.addCallback(future, new FutureCallback<List<EdgeId>>() {
@Override
public void onSuccess(@Nullable List<EdgeId> edgeIds) {
boolean update = false;
if (edgeIds != null && !edgeIds.isEmpty()) {
if (edgeIds.contains(edge.getId())) {
update = true;
}
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<EdgeId> pageData;
do {
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, device.getId(), pageLink);
boolean update = false;
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
if (pageData.getData().contains(edge.getId())) {
update = true;
}
Device device;
if (update) {
log.info("[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " +
"deviceUpdateMsg [{}], Updating device", tenantId, deviceName, edge.getId(), deviceUpdateMsg);
updateDevice(tenantId, edge, deviceUpdateMsg);
} else {
log.info("[{}] Device with name '{}' already exists on the cloud, but not related to this edge [{}]. deviceUpdateMsg [{}]." +
"Creating a new device with random prefix and relate to this edge", tenantId, deviceName, edge.getId(), deviceUpdateMsg);
String newDeviceName = deviceUpdateMsg.getName() + "_" + RandomStringUtils.randomAlphabetic(15);
device = createDevice(tenantId, edge, deviceUpdateMsg, newDeviceName);
ObjectNode body = mapper.createObjectNode();
body.put("conflictName", deviceName);
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ENTITY_MERGE_REQUEST, device.getId(), body);
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
futureToSet.set(null);
}
@Override
public void onFailure(Throwable t) {
log.error("[{}] Failed to get related edge ids by device id [{}], edge [{}]", tenantId, deviceUpdateMsg, edge.getId(), t);
futureToSet.setException(t);
Device newDevice;
if (update) {
log.info("[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " +
"deviceUpdateMsg [{}], Updating device", tenantId, deviceName, edge.getId(), deviceUpdateMsg);
updateDevice(tenantId, edge, deviceUpdateMsg);
} else {
log.info("[{}] Device with name '{}' already exists on the cloud, but not related to this edge [{}]. deviceUpdateMsg [{}]." +
"Creating a new device with random prefix and relate to this edge", tenantId, deviceName, edge.getId(), deviceUpdateMsg);
String newDeviceName = deviceUpdateMsg.getName() + "_" + RandomStringUtils.randomAlphabetic(15);
newDevice = createDevice(tenantId, edge, deviceUpdateMsg, newDeviceName);
ObjectNode body = mapper.createObjectNode();
body.put("conflictName", deviceName);
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ENTITY_MERGE_REQUEST, newDevice.getId(), body);
}
}, dbCallbackExecutorService);
return futureToSet;
} while (pageData != null && pageData.hasNext());
} else {
log.info("[{}] Creating new device and replacing device entity on the edge [{}]", tenantId, deviceUpdateMsg);
device = createDevice(tenantId, edge, deviceUpdateMsg, deviceUpdateMsg.getName());

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java

@ -54,7 +54,7 @@ public class EdgeProcessor extends BaseEdgeProcessor {
public void onSuccess(@Nullable Edge edge) {
if (edge != null && !customerId.isNullUid()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, customerId, null);
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<User> pageData;
do {
pageData = userService.findCustomerUsers(tenantId, customerId, pageLink);

83
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java

@ -93,63 +93,56 @@ public class EntityEdgeProcessor extends BaseEdgeProcessor {
EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type,
new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB()));
EdgeId edgeId = new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB()));
ListenableFuture<List<EdgeId>> edgeIdsFuture;
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<EdgeId> pageData;
switch (actionType) {
case ADDED: // used only for USER entity
case UPDATED:
case CREDENTIALS_UPDATED:
edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId);
Futures.addCallback(edgeIdsFuture, new FutureCallback<List<EdgeId>>() {
@Override
public void onSuccess(@Nullable List<EdgeId> edgeIds) {
if (edgeIds != null && !edgeIds.isEmpty()) {
for (EdgeId edgeId : edgeIds) {
saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null);
}
do {
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (EdgeId relatedEdgeId : pageData.getData()) {
saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
@Override
public void onFailure(Throwable throwable) {
log.error("Failed to find related edge ids [{}]", edgeNotificationMsg, throwable);
}
}, dbCallbackExecutorService);
} while (pageData != null && pageData.hasNext());
break;
case ASSIGNED_TO_CUSTOMER:
case UNASSIGNED_FROM_CUSTOMER:
edgeIdsFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId);
Futures.addCallback(edgeIdsFuture, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable List<EdgeId> edgeIds) {
if (edgeIds != null && !edgeIds.isEmpty()) {
for (EdgeId edgeId : edgeIds) {
try {
CustomerId customerId = mapper.readValue(edgeNotificationMsg.getBody(), CustomerId.class);
ListenableFuture<Edge> future = edgeService.findEdgeByIdAsync(tenantId, edgeId);
Futures.addCallback(future, new FutureCallback<Edge>() {
@Override
public void onSuccess(@Nullable Edge edge) {
if (edge != null && edge.getCustomerId() != null &&
!edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(customerId)) {
saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null);
}
do {
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (EdgeId relatedEdgeId : pageData.getData()) {
try {
CustomerId customerId = mapper.readValue(edgeNotificationMsg.getBody(), CustomerId.class);
ListenableFuture<Edge> future = edgeService.findEdgeByIdAsync(tenantId, relatedEdgeId);
Futures.addCallback(future, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable Edge edge) {
if (edge != null && edge.getCustomerId() != null &&
!edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(customerId)) {
saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null);
}
@Override
public void onFailure(Throwable throwable) {
log.error("Failed to find edge by id [{}]", edgeNotificationMsg, throwable);
}
}, dbCallbackExecutorService);
} catch (Exception e) {
log.error("Can't parse customer id from entity body [{}]", edgeNotificationMsg, e);
}
}
@Override
public void onFailure(Throwable t) {
log.error("Failed to find edge by id [{}] {}", edgeNotificationMsg, t);
}
}, dbCallbackExecutorService);
} catch (Exception e) {
log.error("Can't parse customer id from entity body [{}]", edgeNotificationMsg, e);
}
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
@Override
public void onFailure(Throwable throwable) {
log.error("Failed to find related edge ids [{}]", edgeNotificationMsg, throwable);
}
}, dbCallbackExecutorService);
} while (pageData != null && pageData.hasNext());
break;
case DELETED:
saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null);
@ -165,7 +158,7 @@ public class EntityEdgeProcessor extends BaseEdgeProcessor {
}
private void updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) {
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<RuleChain> pageData;
do {
pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink);

62
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java

@ -16,11 +16,9 @@
package org.thingsboard.server.service.edge.rpc.processor;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.edge.EdgeEvent;
@ -38,6 +36,8 @@ import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.gen.edge.DownlinkMsg;
@ -127,39 +127,35 @@ public class RelationEdgeProcessor extends BaseEdgeProcessor {
EntityRelation relation = mapper.readValue(edgeNotificationMsg.getBody(), EntityRelation.class);
if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) &&
!relation.getTo().getEntityType().equals(EntityType.EDGE)) {
List<ListenableFuture<List<EdgeId>>> futures = new ArrayList<>();
futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo()));
futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom()));
ListenableFuture<List<List<EdgeId>>> combinedFuture = Futures.allAsList(futures);
Futures.addCallback(combinedFuture, new FutureCallback<List<List<EdgeId>>>() {
@Override
public void onSuccess(@Nullable List<List<EdgeId>> listOfListsEdgeIds) {
Set<EdgeId> uniqueEdgeIds = new HashSet<>();
if (listOfListsEdgeIds != null && !listOfListsEdgeIds.isEmpty()) {
for (List<EdgeId> listOfListsEdgeId : listOfListsEdgeIds) {
if (listOfListsEdgeId != null) {
uniqueEdgeIds.addAll(listOfListsEdgeId);
}
}
}
if (!uniqueEdgeIds.isEmpty()) {
for (EdgeId edgeId : uniqueEdgeIds) {
saveEdgeEvent(tenantId,
edgeId,
EdgeEventType.RELATION,
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()),
null,
mapper.valueToTree(relation));
}
}
Set<EdgeId> uniqueEdgeIds = new HashSet<>();
uniqueEdgeIds.addAll(findRelatedEdgeIds(tenantId, relation.getTo()));
uniqueEdgeIds.addAll(findRelatedEdgeIds(tenantId, relation.getFrom()));
if (!uniqueEdgeIds.isEmpty()) {
for (EdgeId edgeId : uniqueEdgeIds) {
saveEdgeEvent(tenantId,
edgeId,
EdgeEventType.RELATION,
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()),
null,
mapper.valueToTree(relation));
}
}
}
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] can't find related edge ids by relation to id [{}] and relation from id [{}]" ,
tenantId.getId(), relation.getTo().getId(), relation.getFrom().getId(), t);
private List<EdgeId> findRelatedEdgeIds(TenantId tenantId, EntityId entityId) {
List<EdgeId> result = new ArrayList<>();
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<EdgeId> pageData;
do {
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
result.addAll(pageData.getData());
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}, dbCallbackExecutorService);
}
}
} while (pageData != null && pageData.hasNext());
return result;
}
}

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java

@ -87,7 +87,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
private static final ObjectMapper mapper = new ObjectMapper();
private static final int DEFAULT_LIMIT = 100;
private static final int DEFAULT_PAGE_SIZE = 1000;
@Autowired
private EdgeEventService edgeEventService;
@ -351,7 +351,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
List<ListenableFuture<EdgeEvent>> futures = new ArrayList<>();
log.trace("[{}] syncDevices [{}][{}]", tenantId, edge.getName(), deviceType);
try {
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<Device> pageData;
do {
pageData = deviceService.findDevicesByTenantIdAndEdgeIdAndType(tenantId, edge.getId(), deviceType, pageLink);

6
common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java

@ -82,11 +82,9 @@ public interface EdgeService {
void assignDefaultRuleChainsToEdge(TenantId tenantId, EdgeId edgeId);
ListenableFuture<List<Edge>> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId);
PageData<Edge> findEdgesByTenantIdAndEntityId(TenantId tenantId, EntityId ruleChainId, PageLink pageLink);
ListenableFuture<List<Edge>> findEdgesByTenantIdAndDashboardId(TenantId tenantId, DashboardId dashboardId);
ListenableFuture<List<EdgeId>> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId);
PageData<EdgeId> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink);
Object checkInstance(Object request);

18
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java

@ -17,6 +17,7 @@ package org.thingsboard.server.dao.edge;
import com.google.common.util.concurrent.ListenableFuture;
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.edge.EdgeInfo;
import org.thingsboard.server.common.data.id.TenantId;
@ -160,20 +161,13 @@ public interface EdgeDao extends Dao<Edge> {
PageData<EdgeInfo> findEdgeInfosByTenantId(UUID tenantId, PageLink pageLink);
/**
* Find edges by tenantId and ruleChainId.
* Find edges by tenantId and entityId.
*
* @param tenantId the tenantId
* @param ruleChainId the ruleChainId
* @return the list of rule chain objects
* @param entityId the entityId
* @param entityType the entityType
* @return the list of edge objects
*/
ListenableFuture<List<Edge>> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId);
PageData<Edge> findEdgesByTenantIdAndEntityId(UUID tenantId, UUID entityId, EntityType entityType, PageLink pageLink);
/**
* Find edges by tenantId and dashboardId.
*
* @param tenantId the tenantId
* @param dashboardId the dashboardId
* @return the list of rule chain objects
*/
ListenableFuture<List<Edge>> findEdgesByTenantIdAndDashboardId(UUID tenantId, UUID dashboardId);
}

117
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java

@ -48,7 +48,6 @@ import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeInfo;
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;
@ -59,7 +58,6 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
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.common.data.rule.RuleChainConnectionInfo;
import org.thingsboard.server.dao.customer.CustomerDao;
@ -106,7 +104,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
private static final ObjectMapper mapper = new ObjectMapper();
private static final int DEFAULT_LIMIT = 100;
private static final int DEFAULT_PAGE_SIZE = 1000;
private RestTemplate restTemplate;
@ -380,7 +378,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
@Override
public void assignDefaultRuleChainsToEdge(TenantId tenantId, EdgeId edgeId) {
log.trace("Executing assignDefaultRuleChainsToEdge, tenantId [{}], edgeId [{}]", tenantId, edgeId);
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<RuleChain> pageData;
do {
pageData = ruleChainService.findAutoAssignToEdgeRuleChainsByTenantId(tenantId, pageLink);
@ -396,19 +394,11 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
}
@Override
public ListenableFuture<List<Edge>> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId) {
log.trace("Executing findEdgesByTenantIdAndRuleChainId, tenantId [{}], ruleChainId [{}]", tenantId, ruleChainId);
public PageData<Edge> findEdgesByTenantIdAndEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) {
log.trace("Executing findEdgesByTenantIdAndEntityId, tenantId [{}], entityId [{}], pageLink [{}]", tenantId, entityId, pageLink);
Validator.validateId(tenantId, "Incorrect tenantId " + tenantId);
Validator.validateId(ruleChainId, "Incorrect ruleChainId " + ruleChainId);
return edgeDao.findEdgesByTenantIdAndRuleChainId(tenantId.getId(), ruleChainId.getId());
}
@Override
public ListenableFuture<List<Edge>> findEdgesByTenantIdAndDashboardId(TenantId tenantId, DashboardId dashboardId) {
log.trace("Executing findEdgesByTenantIdAndDashboardId, tenantId [{}], dashboardId [{}]", tenantId, dashboardId);
Validator.validateId(tenantId, "Incorrect tenantId " + tenantId);
Validator.validateId(dashboardId, "Incorrect dashboardId " + dashboardId);
return edgeDao.findEdgesByTenantIdAndDashboardId(tenantId.getId(), dashboardId.getId());
validatePageLink(pageLink);
return edgeDao.findEdgesByTenantIdAndEntityId(tenantId.getId(), entityId.getId(), entityId.getEntityType(), pageLink);
}
private DataValidator<Edge> edgeValidator =
@ -496,88 +486,55 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
};
@Override
public ListenableFuture<List<EdgeId>> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) {
// TODO: @voba - rewrite 'find' to use native SQL queries instead of fetching relations
log.trace("[{}] Executing findRelatedEdgeIdsByEntityId [{}]", tenantId, entityId);
public PageData<EdgeId> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) {
log.trace("[{}] Executing findRelatedEdgeIdsByEntityId [{}] [{}]", tenantId, entityId, pageLink);
if (EntityType.TENANT.equals(entityId.getEntityType()) ||
EntityType.CUSTOMER.equals(entityId.getEntityType()) ||
EntityType.DEVICE_PROFILE.equals(entityId.getEntityType())) {
List<EdgeId> result = new ArrayList<>();
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageData<Edge> pageData;
do {
if (EntityType.TENANT.equals(entityId.getEntityType()) ||
EntityType.DEVICE_PROFILE.equals(entityId.getEntityType())) {
pageData = findEdgesByTenantId(tenantId, pageLink);
} else {
pageData = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), pageLink);
}
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (Edge edge : pageData.getData()) {
result.add(edge.getId());
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
return Futures.immediateFuture(result);
if (EntityType.TENANT.equals(entityId.getEntityType()) ||
EntityType.DEVICE_PROFILE.equals(entityId.getEntityType())) {
return convertToEdgeIds(findEdgesByTenantId(tenantId, pageLink));
} else {
return convertToEdgeIds(findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), pageLink));
}
} else {
switch (entityId.getEntityType()) {
case DEVICE:
case ASSET:
case ENTITY_VIEW:
ListenableFuture<List<EntityRelation>> originatorEdgeRelationsFuture =
relationService.findByToAndTypeAsync(tenantId, entityId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
return Futures.transform(originatorEdgeRelationsFuture, originatorEdgeRelations -> {
if (originatorEdgeRelations != null && originatorEdgeRelations.size() > 0 &&
originatorEdgeRelations.get(0).getFrom() != null) {
return Collections.singletonList(new EdgeId(originatorEdgeRelations.get(0).getFrom().getId()));
} else {
return Collections.emptyList();
}
}, MoreExecutors.directExecutor());
case DASHBOARD:
return convertToEdgeIds(findEdgesByTenantIdAndDashboardId(tenantId, new DashboardId(entityId.getId())));
case RULE_CHAIN:
return convertToEdgeIds(findEdgesByTenantIdAndRuleChainId(tenantId, new RuleChainId(entityId.getId())));
return convertToEdgeIds(findEdgesByTenantIdAndEntityId(tenantId, entityId, pageLink));
case USER:
User userById = userService.findUserById(tenantId, new UserId(entityId.getId()));
if (userById == null) {
return Futures.immediateFuture(Collections.emptyList());
return createEmptyEdgeIdPageData();
}
if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) {
return convertToEdgeIds(findEdgesByTenantId(tenantId, pageLink));
} else {
return convertToEdgeIds(findEdgesByTenantIdAndCustomerId(tenantId, userById.getCustomerId(), pageLink));
}
List<Edge> result = new ArrayList<>();
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageData<Edge> pageData;
do {
if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) {
pageData = findEdgesByTenantId(tenantId, pageLink);
} else {
pageData = findEdgesByTenantIdAndCustomerId(tenantId, userById.getCustomerId(), pageLink);
}
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
result.addAll(pageData.getData());
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
return convertToEdgeIds(Futures.immediateFuture(result));
default:
return Futures.immediateFuture(Collections.emptyList());
log.warn("[{}] Unsupported entity type {}", tenantId, entityId.getEntityType());
return createEmptyEdgeIdPageData();
}
}
}
private ListenableFuture<List<EdgeId>> convertToEdgeIds(ListenableFuture<List<Edge>> future) {
return Futures.transform(future, edges -> {
if (edges != null && !edges.isEmpty()) {
return edges.stream().map(IdBased::getId).collect(Collectors.toList());
} else {
return Collections.emptyList();
}
}, MoreExecutors.directExecutor());
private PageData<EdgeId> createEmptyEdgeIdPageData() {
return new PageData<>(new ArrayList<>(), 0, 0, false);
}
private PageData<EdgeId> convertToEdgeIds(PageData<Edge> pageData) {
if (pageData == null) {
return createEmptyEdgeIdPageData();
}
List<EdgeId> edgeIds = new ArrayList<>();
if (pageData.getData() != null && !pageData.getData().isEmpty()) {
edgeIds = pageData.getData().stream().map(IdBased::getId).collect(Collectors.toList());
}
return new PageData<>(edgeIds, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext());
}
@Override
@ -625,7 +582,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
private List<RuleChain> findEdgeRuleChains(TenantId tenantId, EdgeId edgeId) {
List<RuleChain> result = new ArrayList<>();
PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<RuleChain> pageData;
do {
pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink);

4
dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java

@ -77,8 +77,8 @@ public abstract class AbstractEntityService {
List<EntityView> entityViews = entityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId).get();
if (entityViews != null && !entityViews.isEmpty()) {
EntityView entityView = entityViews.get(0);
// TODO: @voba - refactor this blocking operation in 3.3+
Boolean relationExists = relationService.checkRelation(tenantId,edgeId, entityView.getId(),
// TODO: @voba - refactor this blocking operation
Boolean relationExists = relationService.checkRelation(tenantId, edgeId, entityView.getId(),
EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE).get();
if (relationExists) {
throw new DataValidationException("Can't unassign device/asset from edge that is related to entity view and entity view is assigned to edge!");

22
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -66,7 +66,6 @@ import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.TENANT;
@ -375,18 +374,21 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
throw new DataValidationException("Deletion of Root Tenant Rule Chain is prohibited!");
}
if (RuleChainType.EDGE.equals(ruleChain.getType())) {
try {
List<Edge> edges = edgeService.findEdgesByTenantIdAndRuleChainId(tenantId, ruleChainId).get();
if (edges != null && !edges.isEmpty()) {
for (Edge edge : edges) {
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<Edge> pageData;
do {
pageData = edgeService.findEdgesByTenantIdAndEntityId(tenantId, ruleChainId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (Edge edge : pageData.getData()) {
if (edge.getRootRuleChainId() != null && edge.getRootRuleChainId().equals(ruleChainId)) {
throw new DataValidationException("Can't delete rule chain that is root for edge [" + edge.getName() + "]. Please assign another root rule chain first to the edge!");
}
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
} catch (InterruptedException | ExecutionException e) {
log.error("Can't get edges by tenant id [{}] and rule chain id [{}]", tenantId.getId(), ruleChainId.getId(), e);
}
} while (pageData != null && pageData.hasNext());
}
}
checkRuleNodesAndDelete(tenantId, ruleChainId);
@ -638,7 +640,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
}
private void checkRuleNodesAndDelete(TenantId tenantId, RuleChainId ruleChainId) {
try{
try {
ruleChainDao.removeById(tenantId, ruleChainId.getId());
} catch (Exception t) {
ConstraintViolationException e = extractConstraintViolationException(t).orElse(null);
@ -673,7 +675,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
@Override
protected void validateCreate(TenantId tenantId, RuleChain data) {
DefaultTenantProfileConfiguration profileConfiguration =
(DefaultTenantProfileConfiguration)tenantProfileCache.get(tenantId).getProfileData().getConfiguration();
(DefaultTenantProfileConfiguration) tenantProfileCache.get(tenantId).getProfileData().getConfiguration();
long maxRuleChains = profileConfiguration.getMaxRuleChains();
validateNumberOfEntitiesPerTenant(tenantId, ruleChainDao, maxRuleChains, EntityType.RULE_CHAIN);
}

11
dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java

@ -22,6 +22,7 @@ import org.springframework.data.repository.PagingAndSortingRepository;
import org.springframework.data.repository.query.Param;
import org.thingsboard.server.dao.model.sql.EdgeEntity;
import org.thingsboard.server.dao.model.sql.EdgeInfoEntity;
import org.thingsboard.server.dao.model.sql.RuleChainEntity;
import java.util.List;
import java.util.UUID;
@ -110,6 +111,16 @@ public interface EdgeRepository extends PagingAndSortingRepository<EdgeEntity, U
@Param("textSearch") String textSearch,
Pageable pageable);
@Query("SELECT ee FROM EdgeEntity ee, RelationEntity re WHERE ee.tenantId = :tenantId " +
"AND ee.id = re.fromId AND re.fromType = 'EDGE' AND re.relationTypeGroup = 'EDGE' " +
"AND re.relationType = 'Contains' AND re.toId = :entityId AND re.toType = :entityType " +
"AND LOWER(ee.searchText) LIKE LOWER(CONCAT(:searchText, '%'))")
Page<EdgeEntity> findByTenantIdAndEntityId(@Param("tenantId") UUID tenantId,
@Param("entityId") UUID entityId,
@Param("entityType") String entityType,
@Param("searchText") String searchText,
Pageable pageable);
@Query("SELECT DISTINCT d.type FROM EdgeEntity d WHERE d.tenantId = :tenantId")
List<String> findTenantEdgeTypes(@Param("tenantId") UUID tenantId);

23
dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java

@ -26,13 +26,10 @@ 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.edge.EdgeInfo;
import org.thingsboard.server.common.data.id.DashboardId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.edge.EdgeDao;
import org.thingsboard.server.dao.model.sql.EdgeEntity;
@ -181,17 +178,15 @@ public class JpaEdgeDao extends JpaAbstractSearchTextDao<EdgeEntity, Edge> imple
}
@Override
public ListenableFuture<List<Edge>> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId) {
log.debug("Try to find edges by tenantId [{}], ruleChainId [{}]", tenantId, ruleChainId);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
return transformFromRelationToEdge(tenantId, relations);
}
@Override
public ListenableFuture<List<Edge>> findEdgesByTenantIdAndDashboardId(UUID tenantId, UUID dashboardId) {
log.debug("Try to find edges by tenantId [{}], dashboardId [{}]", tenantId, dashboardId);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new DashboardId(dashboardId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
return transformFromRelationToEdge(tenantId, relations);
public PageData<Edge> findEdgesByTenantIdAndEntityId(UUID tenantId, UUID entityId, EntityType entityType, PageLink pageLink) {
log.debug("Try to find edges by tenantId [{}], entityId [{}], entityType [{}], pageLink [{}]", tenantId, entityId, entityType, pageLink);
return DaoUtil.toPageData(
edgeRepository.findByTenantIdAndEntityId(
tenantId,
entityId,
entityType.name(),
Objects.toString(pageLink.getTextSearch(), ""),
DaoUtil.toPageable(pageLink)));
}
private ListenableFuture<List<Edge>> transformFromRelationToEdge(UUID tenantId, ListenableFuture<List<EntityRelation>> relations) {

105
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

@ -15,7 +15,6 @@
*/
package org.thingsboard.rule.engine.edge;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.FutureCallback;
@ -39,6 +38,8 @@ import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.rule.RuleChainType;
import org.thingsboard.server.common.msg.TbMsg;
@ -90,6 +91,8 @@ public class TbMsgPushToEdgeNode implements TbNode {
private static final String SCOPE = "scope";
private static final int DEFAULT_PAGE_SIZE = 1000;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, EmptyNodeConfiguration.class);
@ -117,77 +120,53 @@ public class TbMsgPushToEdgeNode implements TbNode {
private void processMsg(TbContext ctx, TbMsg msg) {
if (EntityType.EDGE.equals(msg.getOriginator().getEntityType())) {
try {
EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx);
if (edgeEvent != null) {
EdgeId edgeId = new EdgeId(msg.getOriginator().getId());
edgeEvent.setEdgeId(edgeId);
ListenableFuture<EdgeEvent> saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent);
Futures.addCallback(saveFuture, new FutureCallback<EdgeEvent>() {
@Override
public void onSuccess(@Nullable EdgeEvent event) {
ctx.tellNext(msg, SUCCESS);
ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId);
}
@Override
public void onFailure(Throwable th) {
log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th);
ctx.tellFailure(msg, th);
}
}, ctx.getDbCallbackExecutor());
}
} catch (JsonProcessingException e) {
log.error("Failed to build edge event", e);
ctx.tellFailure(msg, e);
EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx);
if (edgeEvent != null) {
EdgeId edgeId = new EdgeId(msg.getOriginator().getId());
notifyEdge(ctx, msg, edgeEvent, edgeId);
}
} else {
ListenableFuture<List<EdgeId>> getEdgeIdsFuture = ctx.getEdgeService().findRelatedEdgeIdsByEntityId(ctx.getTenantId(), msg.getOriginator());
Futures.addCallback(getEdgeIdsFuture, new FutureCallback<List<EdgeId>>() {
@Override
public void onSuccess(@Nullable List<EdgeId> edgeIds) {
if (edgeIds != null && !edgeIds.isEmpty()) {
for (EdgeId edgeId : edgeIds) {
try {
EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx);
if (edgeEvent == null) {
log.debug("Edge event type is null. Entity Type {}", msg.getOriginator().getEntityType());
ctx.tellFailure(msg, new RuntimeException("Edge event type is null. Entity Type '" + msg.getOriginator().getEntityType() + "'"));
} else {
edgeEvent.setEdgeId(edgeId);
ListenableFuture<EdgeEvent> saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent);
Futures.addCallback(saveFuture, new FutureCallback<EdgeEvent>() {
@Override
public void onSuccess(@Nullable EdgeEvent event) {
ctx.tellNext(msg, SUCCESS);
ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId);
}
@Override
public void onFailure(Throwable th) {
log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th);
ctx.tellFailure(msg, th);
}
}, ctx.getDbCallbackExecutor());
}
} catch (JsonProcessingException e) {
log.error("Failed to build edge event", e);
ctx.tellFailure(msg, e);
}
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<EdgeId> pageData;
do {
pageData = ctx.getEdgeService().findRelatedEdgeIdsByEntityId(ctx.getTenantId(), msg.getOriginator(), pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (EdgeId edgeId : pageData.getData()) {
EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx);
if (edgeEvent == null) {
log.debug("Edge event type is null. Entity Type {}", msg.getOriginator().getEntityType());
ctx.tellFailure(msg, new RuntimeException("Edge event type is null. Entity Type '" + msg.getOriginator().getEntityType() + "'"));
} else {
notifyEdge(ctx, msg, edgeEvent, edgeId);
}
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
}
}
@Override
public void onFailure(Throwable t) {
ctx.tellFailure(msg, t);
}
private void notifyEdge(TbContext ctx, TbMsg msg, EdgeEvent edgeEvent, EdgeId edgeId) {
edgeEvent.setEdgeId(edgeId);
ListenableFuture<EdgeEvent> saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent);
Futures.addCallback(saveFuture, new FutureCallback<EdgeEvent>() {
@Override
public void onSuccess(@Nullable EdgeEvent event) {
ctx.tellNext(msg, SUCCESS);
ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId);
}
}, ctx.getDbCallbackExecutor());
}
@Override
public void onFailure(Throwable th) {
log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th);
ctx.tellFailure(msg, th);
}
}, ctx.getDbCallbackExecutor());
}
private EdgeEvent buildEdgeEvent(TbMsg msg, TbContext ctx) throws JsonProcessingException {
private EdgeEvent buildEdgeEvent(TbMsg msg, TbContext ctx) {
String msgType = msg.getType();
if (DataConstants.ALARM.equals(msgType)) {
return buildEdgeEvent(ctx.getTenantId(), EdgeEventActionType.ADDED, getUUIDFromMsgData(msg), EdgeEventType.ALARM, null);

Loading…
Cancel
Save