Browse Source

Merge remote-tracking branch 'origin/develop/2.6-edge' into develop/3.3-edge

pull/3811/head
Volodymyr Babak 6 years ago
parent
commit
64694532db
  1. 16
      application/src/main/data/upgrade/2.6.0/schema_update.cql
  2. 19
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  3. 20
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  4. 7
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  5. 9
      application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
  6. 1
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  7. 7
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  8. 248
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  9. 91
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  10. 63
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  11. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeRpcService.java
  12. 252
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java
  13. 17
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseProcessor.java
  14. 20
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  15. 7
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  16. 3
      application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java
  17. 1
      application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java
  18. 1
      application/src/main/resources/thingsboard.yml
  19. 27
      application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java
  20. 22
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  21. 3
      common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java
  22. 17
      common/edge-api/src/main/proto/edge.proto
  23. 7
      common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java
  24. 42
      common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeEventUpdateMsg.java
  25. 1
      common/queue/src/main/proto/queue.proto
  26. 67
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  27. 3
      dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java
  28. 16
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
  29. 3
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  30. 101
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java
  31. 14
      ui-ngx/src/assets/locale/locale.constant-en_US.json
  32. 2
      ui/src/app/edge/edge.controller.js
  33. 2
      ui/src/app/edge/edge.routes.js

16
application/src/main/data/upgrade/2.6.0/schema_update.cql

@ -107,15 +107,15 @@ CREATE MATERIALIZED VIEW IF NOT EXISTS thingsboard.edge_by_customer_by_type_and_
WITH CLUSTERING ORDER BY ( tenant_id DESC, type ASC, search_text ASC, id DESC ); WITH CLUSTERING ORDER BY ( tenant_id DESC, type ASC, search_text ASC, id DESC );
CREATE TABLE IF NOT EXISTS thingsboard.edge_event ( CREATE TABLE IF NOT EXISTS thingsboard.edge_event (
id timeuuid, id timeuuid,
tenant_id timeuuid, tenant_id timeuuid,
edge_id timeuuid, edge_id timeuuid,
edge_event_type text, edge_event_type text,
edge_event_action text, edge_event_action text,
edge_event_uid text, edge_event_uid text,
entity_id timeuuid, entity_id timeuuid,
body text, body text,
PRIMARY KEY ((tenant_id, edge_id), edge_event_type, edge_event_uid) PRIMARY KEY ((tenant_id, edge_id), edge_event_type, edge_event_uid)
); );
CREATE MATERIALIZED VIEW IF NOT EXISTS thingsboard.edge_event_by_id AS CREATE MATERIALIZED VIEW IF NOT EXISTS thingsboard.edge_event_by_id AS

19
application/src/main/java/org/thingsboard/server/actors/app/AppActor.java

@ -35,10 +35,12 @@ import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.aware.TenantAwareMsg; import org.thingsboard.server.common.msg.aware.TenantAwareMsg;
import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException; import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper;
@ -95,6 +97,9 @@ public class AppActor extends ContextAwareActor {
case SERVER_RPC_RESPONSE_TO_DEVICE_ACTOR_MSG: case SERVER_RPC_RESPONSE_TO_DEVICE_ACTOR_MSG:
onToDeviceActorMsg((TenantAwareMsg) msg, true); onToDeviceActorMsg((TenantAwareMsg) msg, true);
break; break;
case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG:
onToTenantActorMsg((EdgeEventUpdateMsg) msg);
break;
default: default:
return false; return false;
} }
@ -194,6 +199,20 @@ public class AppActor extends ContextAwareActor {
() -> new TenantActor.ActorCreator(systemContext, tenantId)); () -> new TenantActor.ActorCreator(systemContext, tenantId));
} }
private void onToTenantActorMsg(EdgeEventUpdateMsg msg) {
TbActorRef target = null;
if (ModelConstants.SYSTEM_TENANT.equals(msg.getTenantId())) {
log.warn("Message has system tenant id: {}", msg);
} else {
target = getOrCreateTenantActor(msg.getTenantId());
}
if (target != null) {
target.tellWithHighPriority(msg);
} else {
log.debug("[{}] Invalid edge event update msg: {}", msg.getTenantId(), msg);
}
}
public static class ActorCreator extends ContextBasedCreator { public static class ActorCreator extends ContextBasedCreator {
public ActorCreator(ActorSystemContext context) { public ActorCreator(ActorSystemContext context) {

20
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -145,8 +145,13 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
if (result != null && result.size() > 0) { if (result != null && result.size() > 0) {
EntityRelation relationToEdge = result.get(0); EntityRelation relationToEdge = result.get(0);
if (relationToEdge.getFrom() != null && relationToEdge.getFrom().getId() != null) { if (relationToEdge.getFrom() != null && relationToEdge.getFrom().getId() != null) {
log.trace("[{}][{}] found edge [{}] for device", tenantId, deviceId, relationToEdge.getFrom().getId());
return new EdgeId(relationToEdge.getFrom().getId()); return new EdgeId(relationToEdge.getFrom().getId());
} else {
log.trace("[{}][{}] edge relation is empty {}", tenantId, deviceId, relationToEdge);
} }
} else {
log.trace("[{}][{}] device doesn't have any related edge", tenantId, deviceId);
} }
return null; return null;
} }
@ -165,6 +170,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
boolean sent; boolean sent;
if (systemContext.isEdgesRpcEnabled() && edgeId != null) { if (systemContext.isEdgesRpcEnabled() && edgeId != null) {
log.debug("[{}][{}] device is related to edge [{}]. Saving RPC request to edge queue", tenantId, deviceId, edgeId.getId());
saveRpcRequestToEdgeQueue(request, rpcRequest.getRequestId()); saveRpcRequestToEdgeQueue(request, rpcRequest.getRequestId());
sent = true; sent = true;
} else { } else {
@ -516,6 +522,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
} }
void processEdgeUpdate(DeviceEdgeUpdateMsg msg) { void processEdgeUpdate(DeviceEdgeUpdateMsg msg) {
log.trace("[{}] Processing edge update {}", deviceId, msg);
this.edgeId = msg.getEdgeId(); this.edgeId = msg.getEdgeId();
} }
@ -568,7 +575,18 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
edgeEvent.setBody(body); edgeEvent.setBody(body);
edgeEvent.setEdgeId(edgeId); edgeEvent.setEdgeId(edgeId);
systemContext.getEdgeEventService().saveAsync(edgeEvent); ListenableFuture<EdgeEvent> future = systemContext.getEdgeEventService().saveAsync(edgeEvent);
Futures.addCallback(future, new FutureCallback<EdgeEvent>() {
@Override
public void onSuccess( EdgeEvent result) {
systemContext.getClusterService().onEdgeEventUpdate(tenantId, edgeId);
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] Can't save edge event [{}] for edge [{}]", tenantId.getId(), edgeEvent, edgeId.getId(), t);
}
}, systemContext.getDbCallbackExecutor());
} }
private List<TsKvProto> toTsKvProtos(@Nullable List<AttributeKvEntry> result) { private List<TsKvProto> toTsKvProtos(@Nullable List<AttributeKvEntry> result) {

7
application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java

@ -34,7 +34,6 @@ import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.rule.engine.api.sms.SmsSenderFactory; import org.thingsboard.rule.engine.api.sms.SmsSenderFactory;
import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.TbActorRef; import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
@ -43,6 +42,7 @@ import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.RuleNodeId;
@ -280,6 +280,11 @@ class DefaultTbContext implements TbContext {
return entityActionMsg(alarm, alarm.getId(), ruleNodeId, action); return entityActionMsg(alarm, alarm.getId(), ruleNodeId, action);
} }
@Override
public void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId) {
mainCtx.getClusterService().onEdgeEventUpdate(tenantId, edgeId);
}
public <E, I extends EntityId> TbMsg entityActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action) { public <E, I extends EntityId> TbMsg entityActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action) {
try { try {
return TbMsg.newMsg(action, id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity))); return TbMsg.newMsg(action, id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)));

9
application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java

@ -47,6 +47,7 @@ import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg;
import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg;
import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg;
import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
@ -169,6 +170,9 @@ public class TenantActor extends RuleChainManagerActor {
case RULE_CHAIN_TO_RULE_CHAIN_MSG: case RULE_CHAIN_TO_RULE_CHAIN_MSG:
onRuleChainMsg((RuleChainAwareMsg) msg); onRuleChainMsg((RuleChainAwareMsg) msg);
break; break;
case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG:
onToEdgeSessionMsg((EdgeEventUpdateMsg) msg);
break;
default: default:
return false; return false;
} }
@ -271,6 +275,11 @@ public class TenantActor extends RuleChainManagerActor {
() -> new DeviceActorCreator(systemContext, tenantId, deviceId)); () -> new DeviceActorCreator(systemContext, tenantId, deviceId));
} }
private void onToEdgeSessionMsg(EdgeEventUpdateMsg msg) {
log.trace("[{}] onToEdgeSessionMsg [{}]", msg.getTenantId(), msg);
systemContext.getEdgeRpcService().onEdgeEvent(msg.getEdgeId());
}
public static class ActorCreator extends ContextBasedCreator { public static class ActorCreator extends ContextBasedCreator {
private final TenantId tenantId; private final TenantId tenantId;

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

@ -955,6 +955,7 @@ public abstract class BaseController {
builder.setBody(body); builder.setBody(body);
} }
TransportProtos.EdgeNotificationMsgProto msg = builder.build(); TransportProtos.EdgeNotificationMsgProto msg = builder.build();
log.trace("[{}] sending notification to edge service {}", tenantId.getId(), msg);
tbClusterService.pushMsgToCore(tenantId, entityId != null ? entityId : tenantId, tbClusterService.pushMsgToCore(tenantId, entityId != null ? entityId : tenantId,
TransportProtos.ToCoreMsg.newBuilder().setEdgeNotificationMsg(msg).build(), null); TransportProtos.ToCoreMsg.newBuilder().setEdgeNotificationMsg(msg).build(), null);
} }

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

@ -282,12 +282,11 @@ public class RuleChainController extends BaseController {
try { try {
TenantId tenantId = getCurrentUser().getTenantId(); TenantId tenantId = getCurrentUser().getTenantId();
PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder);
RuleChainType type = RuleChainType.CORE;
if (typeStr != null && typeStr.trim().length() > 0) { if (typeStr != null && typeStr.trim().length() > 0) {
RuleChainType type = RuleChainType.valueOf(typeStr); type = RuleChainType.valueOf(typeStr);
return checkNotNull(ruleChainService.findTenantRuleChainsByType(tenantId, type, pageLink));
} else {
return checkNotNull(ruleChainService.findTenantRuleChainsByType(tenantId, RuleChainType.CORE, pageLink));
} }
return checkNotNull(ruleChainService.findTenantRuleChainsByType(tenantId, type, pageLink));
} catch (Exception e) { } catch (Exception e) {
throw handleException(e); throw handleException(e);
} }

248
application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java

@ -25,6 +25,7 @@ import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable; import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.Alarm;
@ -54,6 +55,7 @@ import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.queue.TbClusterService;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
@ -73,6 +75,8 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
private static final ObjectMapper mapper = new ObjectMapper(); private static final ObjectMapper mapper = new ObjectMapper();
private static final int DEFAULT_LIMIT = 100;
@Autowired @Autowired
private EdgeService edgeService; private EdgeService edgeService;
@ -88,6 +92,9 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
@Autowired @Autowired
private EdgeEventService edgeEventService; private EdgeEventService edgeEventService;
@Autowired
private TbClusterService clusterService;
@Autowired @Autowired
private DbCallbackExecutorService dbCallbackExecutorService; private DbCallbackExecutorService dbCallbackExecutorService;
@ -136,7 +143,19 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
edgeEvent.setEntityId(entityId.getId()); edgeEvent.setEntityId(entityId.getId());
} }
edgeEvent.setBody(body); edgeEvent.setBody(body);
edgeEventService.saveAsync(edgeEvent); ListenableFuture<EdgeEvent> future = edgeEventService.saveAsync(edgeEvent);
Futures.addCallback(future, new FutureCallback<EdgeEvent>() {
@Override
public void onSuccess(@Nullable EdgeEvent result) {
clusterService.onEdgeEventUpdate(tenantId, edgeId);
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] Can't save edge event [{}] for edge [{}]", tenantId.getId(), edgeEvent, edgeId.getId(), t);
}
}, dbCallbackExecutorService);
} }
@Override @Override
@ -194,13 +213,20 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
public void onSuccess(@Nullable Edge edge) { public void onSuccess(@Nullable Edge edge) {
if (edge != null && !customerId.isNullUid()) { if (edge != null && !customerId.isNullUid()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, customerId, null); saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, customerId, null);
PageData<User> pageData = userService.findCustomerUsers(tenantId, customerId, new PageLink(Integer.MAX_VALUE)); PageLink pageLink = new PageLink(DEFAULT_LIMIT);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { PageData<User> pageData;
log.trace("[{}] [{}] user(s) are going to be added to edge.", edge.getId(), pageData.getData().size()); do {
for (User user : pageData.getData()) { pageData = userService.findCustomerUsers(tenantId, customerId, pageLink);
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.USER, EdgeEventActionType.ADDED, user.getId(), null); if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] user(s) are going to be added to edge.", edge.getId(), pageData.getData().size());
for (User user : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.USER, EdgeEventActionType.ADDED, user.getId(), null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
} }
} }
@ -241,12 +267,19 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
case ADDED: case ADDED:
case UPDATED: case UPDATED:
case DELETED: case DELETED:
PageData<Edge> edgesByTenantId = edgeService.findEdgesByTenantId(tenantId, new PageLink(Integer.MAX_VALUE)); PageLink pageLink = new PageLink(DEFAULT_LIMIT);
if (edgesByTenantId != null && edgesByTenantId.getData() != null && !edgesByTenantId.getData().isEmpty()) { PageData<Edge> pageData;
for (Edge edge : edgesByTenantId.getData()) { do {
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null); pageData = edgeService.findEdgesByTenantId(tenantId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (Edge edge : pageData.getData()) {
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
break; break;
} }
} }
@ -255,21 +288,28 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction());
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType());
EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB()));
PageData<Edge> edgesByTenantId = edgeService.findEdgesByTenantId(tenantId, new PageLink(Integer.MAX_VALUE)); PageLink pageLink = new PageLink(DEFAULT_LIMIT);
if (edgesByTenantId != null && edgesByTenantId.getData() != null && !edgesByTenantId.getData().isEmpty()) { PageData<Edge> pageData;
for (Edge edge : edgesByTenantId.getData()) { do {
switch (actionType) { pageData = edgeService.findEdgesByTenantId(tenantId, pageLink);
case UPDATED: if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
if (!edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(entityId)) { for (Edge edge : pageData.getData()) {
switch (actionType) {
case UPDATED:
if (!edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(entityId)) {
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null);
}
break;
case DELETED:
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null); saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null);
} break;
break; }
case DELETED: }
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null); if (pageData.hasNext()) {
break; pageLink = pageLink.nextPageLink();
} }
} }
} } while (pageData != null && pageData.hasNext());
} }
private void processEntity(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { private void processEntity(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) {
@ -336,12 +376,19 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
break; break;
case DELETED: case DELETED:
PageData<Edge> edgesByTenantId = edgeService.findEdgesByTenantId(tenantId, new PageLink(Integer.MAX_VALUE)); PageLink pageLink = new PageLink(DEFAULT_LIMIT);
if (edgesByTenantId != null && edgesByTenantId.getData() != null && !edgesByTenantId.getData().isEmpty()) { PageData<Edge> pageData;
for (Edge edge : edgesByTenantId.getData()) { do {
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null); pageData = edgeService.findEdgesByTenantId(tenantId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (Edge edge : pageData.getData()) {
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
break; break;
case ASSIGNED_TO_EDGE: case ASSIGNED_TO_EDGE:
case UNASSIGNED_FROM_EDGE: case UNASSIGNED_FROM_EDGE:
@ -355,53 +402,74 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
} }
private void updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) { private void updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) {
PageData<RuleChain> pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, new TimePageLink(Integer.MAX_VALUE)); TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { PageData<RuleChain> pageData;
for (RuleChain ruleChain : pageData.getData()) { do {
if (!ruleChain.getId().equals(processingRuleChainId)) { pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink);
List<RuleChainConnectionInfo> connectionInfos = if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
ruleChainService.loadRuleChainMetaData(ruleChain.getTenantId(), ruleChain.getId()).getRuleChainConnections(); for (RuleChain ruleChain : pageData.getData()) {
if (connectionInfos != null && !connectionInfos.isEmpty()) { if (!ruleChain.getId().equals(processingRuleChainId)) {
for (RuleChainConnectionInfo connectionInfo : connectionInfos) { List<RuleChainConnectionInfo> connectionInfos =
if (connectionInfo.getTargetRuleChainId().equals(processingRuleChainId)) { ruleChainService.loadRuleChainMetaData(ruleChain.getTenantId(), ruleChain.getId()).getRuleChainConnections();
saveEdgeEvent(tenantId, if (connectionInfos != null && !connectionInfos.isEmpty()) {
edgeId, for (RuleChainConnectionInfo connectionInfo : connectionInfos) {
EdgeEventType.RULE_CHAIN_METADATA, if (connectionInfo.getTargetRuleChainId().equals(processingRuleChainId)) {
EdgeEventActionType.UPDATED, saveEdgeEvent(tenantId,
ruleChain.getId(), edgeId,
null); EdgeEventType.RULE_CHAIN_METADATA,
EdgeEventActionType.UPDATED,
ruleChain.getId(),
null);
}
} }
} }
} }
} }
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
} }
private void processAlarm(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { private void processAlarm(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) {
AlarmId alarmId = new AlarmId(new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); AlarmId alarmId = new AlarmId(new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB()));
ListenableFuture<Alarm> alarmFuture = alarmService.findAlarmByIdAsync(tenantId, alarmId); ListenableFuture<Alarm> alarmFuture = alarmService.findAlarmByIdAsync(tenantId, alarmId);
Futures.transform(alarmFuture, alarm -> { Futures.addCallback(alarmFuture, new FutureCallback<Alarm>() {
if (alarm != null) { @Override
EdgeEventType type = getEdgeQueueTypeByEntityType(alarm.getOriginator().getEntityType()); public void onSuccess(@Nullable Alarm alarm) {
if (type != null) { if (alarm != null) {
ListenableFuture<List<EdgeId>> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator()); EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType());
Futures.transform(relatedEdgeIdsByEntityIdFuture, relatedEdgeIdsByEntityId -> { if (type != null) {
if (relatedEdgeIdsByEntityId != null) { ListenableFuture<List<EdgeId>> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator());
for (EdgeId edgeId : relatedEdgeIdsByEntityId) { Futures.addCallback(relatedEdgeIdsByEntityIdFuture, new FutureCallback<List<EdgeId>>() {
saveEdgeEvent(tenantId, @Override
edgeId, public void onSuccess(@Nullable List<EdgeId> relatedEdgeIdsByEntityId) {
EdgeEventType.ALARM, if (relatedEdgeIdsByEntityId != null) {
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), for (EdgeId edgeId : relatedEdgeIdsByEntityId) {
alarmId, saveEdgeEvent(tenantId,
null); edgeId,
EdgeEventType.ALARM,
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()),
alarmId,
null);
}
}
} }
}
return null; @Override
}, dbCallbackExecutorService); public void onFailure(Throwable t) {
log.warn("[{}] can't find related edge ids by entity id [{}]", tenantId.getId(), alarm.getOriginator(), t);
}
}, dbCallbackExecutorService);
}
} }
} }
return null;
@Override
public void onFailure(Throwable t) {
log.warn("[{}] can't find alarm by id [{}]", tenantId.getId(), alarmId.getId(), t);
}
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
} }
@ -413,41 +481,35 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo())); futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getTo()));
futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom())); futures.add(edgeService.findRelatedEdgeIdsByEntityId(tenantId, relation.getFrom()));
ListenableFuture<List<List<EdgeId>>> combinedFuture = Futures.allAsList(futures); ListenableFuture<List<List<EdgeId>>> combinedFuture = Futures.allAsList(futures);
Futures.transform(combinedFuture, listOfListsEdgeIds -> { Futures.addCallback(combinedFuture, new FutureCallback<List<List<EdgeId>>>() {
Set<EdgeId> uniqueEdgeIds = new HashSet<>(); @Override
if (listOfListsEdgeIds != null && !listOfListsEdgeIds.isEmpty()) { public void onSuccess(@Nullable List<List<EdgeId>> listOfListsEdgeIds) {
for (List<EdgeId> listOfListsEdgeId : listOfListsEdgeIds) { Set<EdgeId> uniqueEdgeIds = new HashSet<>();
if (listOfListsEdgeId != null) { if (listOfListsEdgeIds != null && !listOfListsEdgeIds.isEmpty()) {
uniqueEdgeIds.addAll(listOfListsEdgeId); for (List<EdgeId> listOfListsEdgeId : listOfListsEdgeIds) {
if (listOfListsEdgeId != null) {
uniqueEdgeIds.addAll(listOfListsEdgeId);
}
} }
} }
} if (!uniqueEdgeIds.isEmpty()) {
if (!uniqueEdgeIds.isEmpty()) { for (EdgeId edgeId : uniqueEdgeIds) {
for (EdgeId edgeId : uniqueEdgeIds) { saveEdgeEvent(tenantId,
saveEdgeEvent(tenantId, edgeId,
edgeId, EdgeEventType.RELATION,
EdgeEventType.RELATION, EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()),
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), null,
null, mapper.valueToTree(relation));
mapper.valueToTree(relation)); }
} }
} }
return null;
}, dbCallbackExecutorService);
}
}
private EdgeEventType getEdgeQueueTypeByEntityType(EntityType entityType) { @Override
switch (entityType) { public void onFailure(Throwable t) {
case DEVICE: log.warn("[{}] can't find related edge ids by relation to id [{}] and relation from id [{}]" ,
return EdgeEventType.DEVICE; tenantId.getId(), relation.getTo().getId(), relation.getFrom().getId(), t);
case ASSET: }
return EdgeEventType.ASSET; }, dbCallbackExecutorService);
case ENTITY_VIEW:
return EdgeEventType.ENTITY_VIEW;
default:
log.debug("Unsupported entity type: [{}]", entityType);
return null;
} }
} }
} }

91
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -26,6 +26,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
@ -48,9 +49,12 @@ import java.io.File;
import java.io.IOException; import java.io.IOException;
import java.util.Collections; import java.util.Collections;
import java.util.Map; import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
@Service @Service
@ -59,7 +63,9 @@ import java.util.concurrent.TimeUnit;
@TbCoreComponent @TbCoreComponent
public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService {
private final Map<EdgeId, EdgeGrpcSession> sessions = new ConcurrentHashMap<>(); private final ConcurrentMap<EdgeId, EdgeGrpcSession> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, Boolean> sessionNewEvents = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, ScheduledFuture<?>> sessionEdgeEventChecks = new ConcurrentHashMap<>();
private static final ObjectMapper mapper = new ObjectMapper(); private static final ObjectMapper mapper = new ObjectMapper();
@Value("${edges.rpc.port}") @Value("${edges.rpc.port}")
@ -75,6 +81,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Value("${edges.rpc.client_max_keep_alive_time_sec}") @Value("${edges.rpc.client_max_keep_alive_time_sec}")
private int clientMaxKeepAliveTimeSec; private int clientMaxKeepAliveTimeSec;
@Value("${edges.scheduler_pool_size}")
private int schedulerPoolSize;
@Autowired @Autowired
private EdgeContextComponent ctx; private EdgeContextComponent ctx;
@ -83,7 +92,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private Server server; private Server server;
private ExecutorService executor; private ScheduledExecutorService scheduler;
@PostConstruct @PostConstruct
public void init() { public void init() {
@ -109,9 +118,8 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
log.error("Failed to start Edge RPC server!", e); log.error("Failed to start Edge RPC server!", e);
throw new RuntimeException("Failed to start Edge RPC server!"); throw new RuntimeException("Failed to start Edge RPC server!");
} }
this.scheduler = Executors.newScheduledThreadPool(schedulerPoolSize, ThingsBoardThreadFactory.forName("edge-scheduler"));
log.info("Edge RPC service initialized!"); log.info("Edge RPC service initialized!");
executor = Executors.newSingleThreadExecutor();
processHandleMessages();
} }
@PreDestroy @PreDestroy
@ -119,8 +127,16 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
if (server != null) { if (server != null) {
server.shutdownNow(); server.shutdownNow();
} }
if (executor != null) { for (Map.Entry<EdgeId, ScheduledFuture<?>> entry : sessionEdgeEventChecks.entrySet()) {
executor.shutdownNow(); EdgeId edgeId = entry.getKey();
ScheduledFuture<?> sessionEdgeEventCheck = entry.getValue();
if (sessionEdgeEventCheck != null && !sessionEdgeEventCheck.isCancelled() && !sessionEdgeEventCheck.isDone()) {
sessionEdgeEventCheck.cancel(true);
sessionEdgeEventChecks.remove(edgeId);
}
}
if (scheduler != null) {
scheduler.shutdownNow();
} }
} }
@ -147,14 +163,27 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
log.debug("Closing and removing session for edge [{}]", edgeId); log.debug("Closing and removing session for edge [{}]", edgeId);
session.close(); session.close();
sessions.remove(edgeId); sessions.remove(edgeId);
sessionNewEvents.remove(edgeId);
cancelScheduleEdgeEventsCheck(edgeId);
}
}
@Override
public void onEdgeEvent(EdgeId edgeId) {
log.trace("[{}] onEdgeEvent", edgeId.getId());
if (!sessionNewEvents.get(edgeId)) {
log.trace("[{}] set session new events flag to true", edgeId.getId());
sessionNewEvents.put(edgeId, true);
} }
} }
private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) {
log.debug("[{}] onEdgeConnect [{}]", edgeId, edgeGrpcSession.getSessionId()); log.debug("[{}] onEdgeConnect [{}]", edgeId, edgeGrpcSession.getSessionId());
sessions.put(edgeId, edgeGrpcSession); sessions.put(edgeId, edgeGrpcSession);
sessionNewEvents.put(edgeId, false);
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true); save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true);
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis()); save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis());
scheduleEdgeEventsCheck(edgeGrpcSession);
} }
public EdgeGrpcSession getEdgeGrpcSessionById(EdgeId edgeId) { public EdgeGrpcSession getEdgeGrpcSessionById(EdgeId edgeId) {
@ -166,37 +195,47 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
} }
} }
private void processHandleMessages() { private void scheduleEdgeEventsCheck(EdgeGrpcSession session) {
executor.submit(() -> { EdgeId edgeId = session.getEdge().getId();
while (!Thread.interrupted()) { UUID tenantId = session.getEdge().getTenantId().getId();
if (sessions.containsKey(edgeId)) {
ScheduledFuture<?> schedule = scheduler.schedule(() -> {
try { try {
if (sessions.size() > 0) { if (sessionNewEvents.get(edgeId)) {
for (EdgeGrpcSession session : sessions.values()) { log.trace("[{}] Set session new events flag to false", edgeId.getId());
session.processHandleMessages(); sessionNewEvents.put(edgeId, false);
} session.processEdgeEvents();
} else {
log.trace("No sessions available, sleep for the next run");
try {
Thread.sleep(1000);
} catch (InterruptedException ignore) {
}
} }
} catch (Exception e) { } catch (Exception e) {
log.warn("Failed to process messages handling!", e); log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), e);
try {
Thread.sleep(1000);
} catch (InterruptedException ignore) {
}
} }
scheduleEdgeEventsCheck(session);
}, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval(), TimeUnit.MILLISECONDS);
sessionEdgeEventChecks.put(edgeId, schedule);
log.trace("[{}] Check edge event was scheduler for edge [{}]", tenantId, edgeId.getId());
} else {
log.debug("[{}] Session was removed and edge event check schedule must not be started [{}]",
tenantId, edgeId.getId());
}
}
private void cancelScheduleEdgeEventsCheck(EdgeId edgeId) {
if (sessionEdgeEventChecks.containsKey(edgeId)) {
ScheduledFuture<?> sessionEdgeEventCheck = sessionEdgeEventChecks.get(edgeId);
if (sessionEdgeEventCheck != null && !sessionEdgeEventCheck.isCancelled() && !sessionEdgeEventCheck.isDone()) {
sessionEdgeEventCheck.cancel(true);
sessionEdgeEventChecks.remove(edgeId);
} }
}); }
} }
private void onEdgeDisconnect(EdgeId edgeId) { private void onEdgeDisconnect(EdgeId edgeId) {
log.debug("[{}] onEdgeDisconnect", edgeId); log.debug("[{}] onEdgeDisconnect", edgeId);
sessions.remove(edgeId); sessions.remove(edgeId);
sessionNewEvents.remove(edgeId);
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false); save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false);
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis()); save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis());
cancelScheduleEdgeEventsCheck(edgeId);
} }
private void save(EdgeId edgeId, String key, long value) { private void save(EdgeId edgeId, String key, long value) {

63
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -68,6 +68,7 @@ import org.thingsboard.server.common.data.security.UserCredentials;
import org.thingsboard.server.common.data.widget.WidgetType; import org.thingsboard.server.common.data.widget.WidgetType;
import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.common.data.widget.WidgetsBundle;
import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.common.transport.util.JsonUtils;
import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import org.thingsboard.server.gen.edge.AdminSettingsUpdateMsg; import org.thingsboard.server.gen.edge.AdminSettingsUpdateMsg;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg;
@ -121,7 +122,7 @@ import java.util.function.Consumer;
@Data @Data
public final class EdgeGrpcSession implements Closeable { public final class EdgeGrpcSession implements Closeable {
private static final ReentrantLock responseMsgLock = new ReentrantLock(); private static final ReentrantLock downlinkMsgLock = new ReentrantLock();
private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs";
@ -200,7 +201,7 @@ public final class EdgeGrpcSession implements Closeable {
@Override @Override
public void onSuccess(@Nullable List<Void> result) { public void onSuccess(@Nullable List<Void> result) {
UplinkResponseMsg uplinkResponseMsg = UplinkResponseMsg.newBuilder().setSuccess(true).build(); UplinkResponseMsg uplinkResponseMsg = UplinkResponseMsg.newBuilder().setSuccess(true).build();
sendResponseMsg(ResponseMsg.newBuilder() sendDownlinkMsg(ResponseMsg.newBuilder()
.setUplinkResponseMsg(uplinkResponseMsg) .setUplinkResponseMsg(uplinkResponseMsg)
.build()); .build());
} }
@ -208,7 +209,7 @@ public final class EdgeGrpcSession implements Closeable {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
UplinkResponseMsg uplinkResponseMsg = UplinkResponseMsg.newBuilder().setSuccess(false).setErrorMsg(t.getMessage()).build(); UplinkResponseMsg uplinkResponseMsg = UplinkResponseMsg.newBuilder().setSuccess(false).setErrorMsg(t.getMessage()).build();
sendResponseMsg(ResponseMsg.newBuilder() sendDownlinkMsg(ResponseMsg.newBuilder()
.setUplinkResponseMsg(uplinkResponseMsg) .setUplinkResponseMsg(uplinkResponseMsg)
.build()); .build());
} }
@ -228,38 +229,35 @@ public final class EdgeGrpcSession implements Closeable {
} }
} }
private void sendResponseMsg(ResponseMsg responseMsg) { private void sendDownlinkMsg(ResponseMsg downlinkMsg) {
log.trace("[{}] Sending response msg [{}]", this.sessionId, responseMsg); log.trace("[{}] Sending downlink msg [{}]", this.sessionId, downlinkMsg);
if (isConnected()) { if (isConnected()) {
try { try {
responseMsgLock.lock(); downlinkMsgLock.lock();
outputStream.onNext(responseMsg); outputStream.onNext(downlinkMsg);
} catch (Exception e) { } catch (Exception e) {
log.error("[{}] Failed to send response message [{}]", this.sessionId, responseMsg, e); log.error("[{}] Failed to send downlink message [{}]", this.sessionId, downlinkMsg, e);
connected = false; connected = false;
sessionCloseListener.accept(edge.getId()); sessionCloseListener.accept(edge.getId());
} finally { } finally {
responseMsgLock.unlock(); downlinkMsgLock.unlock();
} }
log.trace("[{}] Response msg successfully sent [{}]", this.sessionId, responseMsg); log.trace("[{}] Response msg successfully sent [{}]", this.sessionId, downlinkMsg);
} }
} }
void onConfigurationUpdate(Edge edge) { void onConfigurationUpdate(Edge edge) {
log.debug("[{}] onConfigurationUpdate [{}]", this.sessionId, edge); log.debug("[{}] onConfigurationUpdate [{}]", this.sessionId, edge);
try { this.edge = edge;
this.edge = edge; EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder()
EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder() .setConfiguration(constructEdgeConfigProto(edge)).build();
.setConfiguration(constructEdgeConfigProto(edge)).build(); ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder()
outputStream.onNext(ResponseMsg.newBuilder() .setEdgeUpdateMsg(edgeConfig)
.setEdgeUpdateMsg(edgeConfig) .build();
.build()); sendDownlinkMsg(edgeConfigMsg);
} catch (Exception e) {
log.error("[{}] Failed to construct proto objects!", this.sessionId, e);
}
} }
void processHandleMessages() throws ExecutionException, InterruptedException { void processEdgeEvents() throws ExecutionException, InterruptedException {
log.trace("[{}] processHandleMessages started", this.sessionId); log.trace("[{}] processHandleMessages started", this.sessionId);
if (isConnected()) { if (isConnected()) {
Long queueStartTs = getQueueStartTs().get(); Long queueStartTs = getQueueStartTs().get();
@ -282,7 +280,7 @@ public final class EdgeGrpcSession implements Closeable {
latch = new CountDownLatch(downlinkMsgsPack.size()); latch = new CountDownLatch(downlinkMsgsPack.size());
for (DownlinkMsg downlinkMsg : downlinkMsgsPack) { for (DownlinkMsg downlinkMsg : downlinkMsgsPack) {
sendResponseMsg(ResponseMsg.newBuilder() sendDownlinkMsg(ResponseMsg.newBuilder()
.setDownlinkMsg(downlinkMsg) .setDownlinkMsg(downlinkMsg)
.build()); .build());
} }
@ -310,11 +308,6 @@ public final class EdgeGrpcSession implements Closeable {
Long newStartTs = Uuids.unixTimestamp(ifOffset); Long newStartTs = Uuids.unixTimestamp(ifOffset);
updateQueueStartTs(newStartTs); updateQueueStartTs(newStartTs);
} }
try {
Thread.sleep(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval());
} catch (InterruptedException e) {
log.error("[{}] Error during sleep between no records interval", this.sessionId, e);
}
} }
log.trace("[{}] processHandleMessages finished", this.sessionId); log.trace("[{}] processHandleMessages finished", this.sessionId);
} }
@ -446,6 +439,9 @@ public final class EdgeGrpcSession implements Closeable {
case CUSTOMER: case CUSTOMER:
entityId = new CustomerId(edgeEvent.getEntityId()); entityId = new CustomerId(edgeEvent.getEntityId());
break; break;
case EDGE:
entityId = new EdgeId(edgeEvent.getEntityId());
break;
} }
DownlinkMsg downlinkMsg = null; DownlinkMsg downlinkMsg = null;
if (entityId != null) { if (entityId != null) {
@ -969,17 +965,24 @@ public final class EdgeGrpcSession implements Closeable {
} }
private EdgeConfiguration constructEdgeConfigProto(Edge edge) { private EdgeConfiguration constructEdgeConfigProto(Edge edge) {
return EdgeConfiguration.newBuilder() EdgeConfiguration.Builder builder = EdgeConfiguration.newBuilder()
.setEdgeIdMSB(edge.getId().getId().getMostSignificantBits()) .setEdgeIdMSB(edge.getId().getId().getMostSignificantBits())
.setEdgeIdLSB(edge.getId().getId().getLeastSignificantBits()) .setEdgeIdLSB(edge.getId().getId().getLeastSignificantBits())
.setTenantIdMSB(edge.getTenantId().getId().getMostSignificantBits()) .setTenantIdMSB(edge.getTenantId().getId().getMostSignificantBits())
.setTenantIdLSB(edge.getTenantId().getId().getLeastSignificantBits()) .setTenantIdLSB(edge.getTenantId().getId().getLeastSignificantBits())
.setName(edge.getName()) .setName(edge.getName())
.setRoutingKey(edge.getRoutingKey())
.setType(edge.getType()) .setType(edge.getType())
.setRoutingKey(edge.getRoutingKey())
.setSecret(edge.getSecret())
.setEdgeLicenseKey(edge.getEdgeLicenseKey()) .setEdgeLicenseKey(edge.getEdgeLicenseKey())
.setCloudEndpoint(edge.getCloudEndpoint()) .setCloudEndpoint(edge.getCloudEndpoint())
.setCloudType("CE") .setConfiguration(JacksonUtil.toString(edge.getConfiguration()))
.setCloudType("CE");
if (edge.getCustomerId() != null) {
builder.setCustomerIdMSB(edge.getCustomerId().getId().getMostSignificantBits())
.setCustomerIdLSB(edge.getCustomerId().getId().getLeastSignificantBits());
}
return builder
.build(); .build();
} }

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

@ -23,4 +23,6 @@ public interface EdgeRpcService {
void updateEdge(Edge edge); void updateEdge(Edge edge);
void deleteEdge(EdgeId edgeId); void deleteEdge(EdgeId edgeId);
void onEdgeEvent(EdgeId edgeId);
} }

252
application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java

@ -35,6 +35,7 @@ import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.DashboardInfo;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.User;
@ -81,6 +82,7 @@ import org.thingsboard.server.gen.edge.RelationRequestMsg;
import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg;
import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg;
import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.queue.TbClusterService;
import java.io.File; import java.io.File;
import java.nio.charset.StandardCharsets; import java.nio.charset.StandardCharsets;
@ -98,6 +100,8 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private static final ObjectMapper mapper = new ObjectMapper(); private static final ObjectMapper mapper = new ObjectMapper();
private static final int DEFAULT_LIMIT = 100;
@Autowired @Autowired
private EdgeEventService edgeEventService; private EdgeEventService edgeEventService;
@ -137,9 +141,12 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
@Autowired @Autowired
private DbCallbackExecutorService dbCallbackExecutorService; private DbCallbackExecutorService dbCallbackExecutorService;
@Autowired
private TbClusterService tbClusterService;
@Override @Override
public void sync(Edge edge) { public void sync(Edge edge) {
log.trace("[{}] staring sync process for edge [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}][{}] Staring edge sync process", edge.getTenantId(), edge.getId());
try { try {
syncWidgetsBundleAndWidgetTypes(edge); syncWidgetsBundleAndWidgetTypes(edge);
syncAdminSettings(edge); syncAdminSettings(edge);
@ -150,21 +157,27 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
syncEntityViews(edge); syncEntityViews(edge);
syncDashboards(edge); syncDashboards(edge);
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during sync process", e); log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), e);
} }
} }
private void syncRuleChains(Edge edge) { private void syncRuleChains(Edge edge) {
log.trace("[{}] syncRuleChains [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncRuleChains [{}]", edge.getTenantId(), edge.getName());
try { try {
PageData<RuleChain> pageData = TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
ruleChainService.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); PageData<RuleChain> pageData;
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { do {
log.trace("[{}] [{}] rule chains(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
for (RuleChain ruleChain : pageData.getData()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.RULE_CHAIN, EdgeEventActionType.ADDED, ruleChain.getId(), null); log.trace("[{}] [{}] rule chains(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size());
for (RuleChain ruleChain : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.RULE_CHAIN, EdgeEventActionType.ADDED, ruleChain.getId(), null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading edge rule chain(s) on sync!", e); log.error("Exception during loading edge rule chain(s) on sync!", e);
} }
@ -173,14 +186,20 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private void syncDevices(Edge edge) { private void syncDevices(Edge edge) {
log.trace("[{}] syncDevices [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncDevices [{}]", edge.getTenantId(), edge.getName());
try { try {
PageData<Device> pageData = TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
deviceService.findDevicesByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); PageData<Device> pageData;
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { do {
log.trace("[{}] [{}] device(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); pageData = deviceService.findDevicesByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
for (Device device : pageData.getData()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ADDED, device.getId(), null); log.trace("[{}] [{}] device(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size());
for (Device device : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ADDED, device.getId(), null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading edge device(s) on sync!", e); log.error("Exception during loading edge device(s) on sync!", e);
} }
@ -189,13 +208,20 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private void syncAssets(Edge edge) { private void syncAssets(Edge edge) {
log.trace("[{}] syncAssets [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncAssets [{}]", edge.getTenantId(), edge.getName());
try { try {
PageData<Asset> pageData = assetService.findAssetsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { PageData<Asset> pageData;
log.trace("[{}] [{}] asset(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); do {
for (Asset asset : pageData.getData()) { pageData = assetService.findAssetsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ASSET, EdgeEventActionType.ADDED, asset.getId(), null); 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()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ASSET, EdgeEventActionType.ADDED, asset.getId(), null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading edge asset(s) on sync!", e); log.error("Exception during loading edge asset(s) on sync!", e);
} }
@ -204,13 +230,20 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private void syncEntityViews(Edge edge) { private void syncEntityViews(Edge edge) {
log.trace("[{}] syncEntityViews [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncEntityViews [{}]", edge.getTenantId(), edge.getName());
try { try {
PageData<EntityView> pageData = entityViewService.findEntityViewsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { PageData<EntityView> pageData;
log.trace("[{}] [{}] entity view(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); do {
for (EntityView entityView : pageData.getData()) { pageData = entityViewService.findEntityViewsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ENTITY_VIEW, EdgeEventActionType.ADDED, entityView.getId(), null); 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()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ENTITY_VIEW, EdgeEventActionType.ADDED, entityView.getId(), null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading edge entity view(s) on sync!", e); log.error("Exception during loading edge entity view(s) on sync!", e);
} }
@ -219,13 +252,20 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private void syncDashboards(Edge edge) { private void syncDashboards(Edge edge) {
log.trace("[{}] syncDashboards [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncDashboards [{}]", edge.getTenantId(), edge.getName());
try { try {
PageData<DashboardInfo> pageData = dashboardService.findDashboardsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE)); TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { PageData<DashboardInfo> pageData;
log.trace("[{}] [{}] dashboard(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); do {
for (DashboardInfo dashboardInfo : pageData.getData()) { pageData = dashboardService.findDashboardsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DASHBOARD, EdgeEventActionType.ADDED, dashboardInfo.getId(), null); 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()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DASHBOARD, EdgeEventActionType.ADDED, dashboardInfo.getId(), null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} }
} } while (pageData != null && pageData.hasNext());
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading edge dashboard(s) on sync!", e); log.error("Exception during loading edge dashboard(s) on sync!", e);
} }
@ -234,18 +274,36 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private void syncUsers(Edge edge) { private void syncUsers(Edge edge) {
log.trace("[{}] syncUsers [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncUsers [{}]", edge.getTenantId(), edge.getName());
try { try {
PageData<User> pageData = userService.findTenantAdmins(edge.getTenantId(), new PageLink(Integer.MAX_VALUE)); TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
pushUsersToEdge(pageData, edge); PageData<User> pageData;
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { do {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null); pageData = userService.findTenantAdmins(edge.getTenantId(), pageLink);
pageData = userService.findCustomerUsers(edge.getTenantId(), edge.getCustomerId(), new PageLink(Integer.MAX_VALUE));
pushUsersToEdge(pageData, edge); pushUsersToEdge(pageData, edge);
} if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} while (pageData.hasNext());
syncCustomerUsers(edge);
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading edge user(s) on sync!", e); log.error("Exception during loading edge user(s) on sync!", e);
} }
} }
private void syncCustomerUsers(Edge edge) {
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null);
TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
PageData<User> pageData;
do {
pageData = userService.findCustomerUsers(edge.getTenantId(), edge.getCustomerId(), pageLink);
pushUsersToEdge(pageData, edge);
if (pageData != null && pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
} while (pageData != null && pageData.hasNext());
}
}
private void syncWidgetsBundleAndWidgetTypes(Edge edge) { private void syncWidgetsBundleAndWidgetTypes(Edge edge) {
log.trace("[{}] syncWidgetsBundleAndWidgetTypes [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncWidgetsBundleAndWidgetTypes [{}]", edge.getTenantId(), edge.getName());
List<WidgetsBundle> widgetsBundlesToPush = new ArrayList<>(); List<WidgetsBundle> widgetsBundlesToPush = new ArrayList<>();
@ -372,10 +430,11 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
EntityId entityId = EntityIdFactory.getByTypeAndUuid( EntityId entityId = EntityIdFactory.getByTypeAndUuid(
EntityType.valueOf(attributesRequestMsg.getEntityType()), EntityType.valueOf(attributesRequestMsg.getEntityType()),
new UUID(attributesRequestMsg.getEntityIdMSB(), attributesRequestMsg.getEntityIdLSB())); new UUID(attributesRequestMsg.getEntityIdMSB(), attributesRequestMsg.getEntityIdLSB()));
final EdgeEventType type = getEdgeQueueTypeByEntityType(entityId.getEntityType()); final EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(entityId.getEntityType());
if (type != null) { if (type != null) {
SettableFuture<Void> futureToSet = SettableFuture.create(); SettableFuture<Void> futureToSet = SettableFuture.create();
ListenableFuture<List<AttributeKvEntry>> ssAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SERVER_SCOPE); String scope = attributesRequestMsg.getScope();
ListenableFuture<List<AttributeKvEntry>> ssAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, scope);
Futures.addCallback(ssAttrFuture, new FutureCallback<List<AttributeKvEntry>>() { Futures.addCallback(ssAttrFuture, new FutureCallback<List<AttributeKvEntry>>() {
@Override @Override
public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) { public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) {
@ -395,7 +454,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
entityData.put("kv", attributes); entityData.put("kv", attributes);
entityData.put("scope", DataConstants.SERVER_SCOPE); entityData.put("scope", scope);
JsonNode body = mapper.valueToTree(entityData); JsonNode body = mapper.valueToTree(entityData);
log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body); log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body);
saveEdgeEvent(edge.getTenantId(), saveEdgeEvent(edge.getTenantId(),
@ -408,6 +467,11 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e); log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e);
throw new RuntimeException("[" + edge.getName() + "] Failed to send attribute updates to the edge", e); throw new RuntimeException("[" + edge.getName() + "] Failed to send attribute updates to the edge", e);
} }
} else {
log.trace("[{}][{}] No attributes found for entity {} [{}]", edge.getTenantId(),
edge.getName(),
entityId.getEntityType(),
entityId.getId());
} }
futureToSet.set(null); futureToSet.set(null);
} }
@ -419,27 +483,12 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
return futureToSet; return futureToSet;
// TODO: voba - push shared attributes to edge?
// ListenableFuture<List<AttributeKvEntry>> shAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SHARED_SCOPE);
// ListenableFuture<List<AttributeKvEntry>> clAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.CLIENT_SCOPE);
} else { } else {
log.warn("[{}] Type doesn't supported {}", edge.getTenantId(), entityId.getEntityType());
return Futures.immediateFuture(null); return Futures.immediateFuture(null);
} }
} }
private EdgeEventType getEdgeQueueTypeByEntityType(EntityType entityType) {
switch (entityType) {
case DEVICE:
return EdgeEventType.DEVICE;
case ASSET:
return EdgeEventType.ASSET;
case ENTITY_VIEW:
return EdgeEventType.ENTITY_VIEW;
default:
return null;
}
}
@Override @Override
public ListenableFuture<Void> processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg) { public ListenableFuture<Void> processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg) {
log.trace("[{}] processRelationRequestMsg [{}][{}]", edge.getTenantId(), edge.getName(), relationRequestMsg); log.trace("[{}] processRelationRequestMsg [{}][{}]", edge.getTenantId(), edge.getName(), relationRequestMsg);
@ -451,34 +500,47 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.FROM)); futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.FROM));
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.TO)); futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.TO));
ListenableFuture<List<List<EntityRelation>>> relationsListFuture = Futures.allAsList(futures); ListenableFuture<List<List<EntityRelation>>> relationsListFuture = Futures.allAsList(futures);
return Futures.transform(relationsListFuture, relationsList -> { SettableFuture<Void> futureToSet = SettableFuture.create();
try { Futures.addCallback(relationsListFuture, new FutureCallback<List<List<EntityRelation>>>() {
if (relationsList != null && !relationsList.isEmpty()) { @Override
for (List<EntityRelation> entityRelations : relationsList) { public void onSuccess(@Nullable List<List<EntityRelation>> relationsList) {
log.trace("[{}] [{}] [{}] relation(s) are going to be pushed to edge.", edge.getId(), entityId, entityRelations.size()); try {
for (EntityRelation relation : entityRelations) { if (relationsList != null && !relationsList.isEmpty()) {
try { for (List<EntityRelation> entityRelations : relationsList) {
if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && log.trace("[{}] [{}] [{}] relation(s) are going to be pushed to edge.", edge.getId(), entityId, entityRelations.size());
!relation.getTo().getEntityType().equals(EntityType.EDGE)) { for (EntityRelation relation : entityRelations) {
saveEdgeEvent(edge.getTenantId(), try {
edge.getId(), if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) &&
EdgeEventType.RELATION, !relation.getTo().getEntityType().equals(EntityType.EDGE)) {
EdgeEventActionType.ADDED, saveEdgeEvent(edge.getTenantId(),
null, edge.getId(),
mapper.valueToTree(relation)); EdgeEventType.RELATION,
EdgeEventActionType.ADDED,
null,
mapper.valueToTree(relation));
}
} catch (Exception e) {
log.error("Exception during loading relation [{}] to edge on sync!", relation, e);
futureToSet.setException(e);
return;
}
}
} }
} catch (Exception e) {
log.error("Exception during loading relation [{}] to edge on sync!", relation, e);
} }
futureToSet.set(null);
} catch (Exception e) {
log.error("Exception during loading relation(s) to edge on sync!", e);
futureToSet.setException(e);
} }
} }
}
} catch (Exception e) { @Override
log.error("Exception during loading relation(s) to edge on sync!", e); public void onFailure(Throwable t) {
throw new RuntimeException("Exception during loading relation(s) to edge on sync!", e); log.error("[{}] Can't find relation by query. Entity id [{}]", edge.getTenantId(), entityId, t);
} futureToSet.setException(t);
return null; }
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
return futureToSet;
} }
private ListenableFuture<List<EntityRelation>> findRelationByQuery(Edge edge, EntityId entityId, EntitySearchDirection direction) { private ListenableFuture<List<EntityRelation>> findRelationByQuery(Edge edge, EntityId entityId, EntitySearchDirection direction) {
@ -534,11 +596,11 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
private ListenableFuture<EdgeEvent> saveEdgeEvent(TenantId tenantId, private ListenableFuture<EdgeEvent> saveEdgeEvent(TenantId tenantId,
EdgeId edgeId, EdgeId edgeId,
EdgeEventType type, EdgeEventType type,
EdgeEventActionType action, EdgeEventActionType action,
EntityId entityId, EntityId entityId,
JsonNode body) { JsonNode body) {
log.trace("Pushing edge event to edge queue. tenantId [{}], edgeId [{}], type [{}], action[{}], entityId [{}], body [{}]", log.trace("Pushing edge event to edge queue. tenantId [{}], edgeId [{}], type [{}], action[{}], entityId [{}], body [{}]",
tenantId, edgeId, type, action, entityId, body); tenantId, edgeId, type, action, entityId, body);
@ -551,6 +613,18 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
edgeEvent.setEntityId(entityId.getId()); edgeEvent.setEntityId(entityId.getId());
} }
edgeEvent.setBody(body); edgeEvent.setBody(body);
return edgeEventService.saveAsync(edgeEvent); ListenableFuture<EdgeEvent> future = edgeEventService.saveAsync(edgeEvent);
Futures.addCallback(future, new FutureCallback<EdgeEvent>() {
@Override
public void onSuccess(@Nullable EdgeEvent result) {
tbClusterService.onEdgeEventUpdate(tenantId, edgeId);
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] Can't save edge event [{}] for edge [{}]", tenantId.getId(), edgeEvent, edgeId.getId(), t);
}
}, dbCallbackExecutorService);
return future;
} }
} }

17
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseProcessor.java

@ -17,8 +17,11 @@ package org.thingsboard.server.service.edge.rpc.processor;
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
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.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType;
@ -107,6 +110,18 @@ public abstract class BaseProcessor {
edgeEvent.setEntityId(entityId.getId()); edgeEvent.setEntityId(entityId.getId());
} }
edgeEvent.setBody(body); edgeEvent.setBody(body);
return edgeEventService.saveAsync(edgeEvent); ListenableFuture<EdgeEvent> future = edgeEventService.saveAsync(edgeEvent);
Futures.addCallback(future, new FutureCallback<EdgeEvent>() {
@Override
public void onSuccess(@Nullable EdgeEvent result) {
tbClusterService.onEdgeEventUpdate(tenantId, edgeId);
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] Can't save edge event [{}] for edge [{}]", tenantId.getId(), edgeEvent, edgeId.getId(), t);
}
}, dbCallbackExecutorService);
return future;
} }
} }

20
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -30,11 +30,13 @@ import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId; 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.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
@ -279,6 +281,21 @@ public class DefaultTbClusterService implements TbClusterService {
} }
} }
@Override
public void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId) {
log.trace("[{}] Processing edge {} event update ", tenantId, edgeId);
EdgeEventUpdateMsg msg = new EdgeEventUpdateMsg(tenantId, edgeId);
byte[] msgBytes = encodingService.encode(msg);
TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer();
Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE);
for (String serviceId : tbCoreServices) {
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceId);
ToCoreNotificationMsg toCoreMsg = ToCoreNotificationMsg.newBuilder().setEdgeEventUpdateMsg(ByteString.copyFrom(msgBytes)).build();
toCoreNfProducer.send(tpi, new TbProtoQueueMsg<>(msg.getEdgeId().getId(), toCoreMsg), null);
toCoreNfs.incrementAndGet();
}
}
private void broadcast(ComponentLifecycleMsg msg) { private void broadcast(ComponentLifecycleMsg msg) {
byte[] msgBytes = encodingService.encode(msg); byte[] msgBytes = encodingService.encode(msg);
TbQueueProducer<TbProtoQueueMsg<ToRuleEngineNotificationMsg>> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToRuleEngineNotificationMsg>> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer();
@ -286,7 +303,8 @@ public class DefaultTbClusterService implements TbClusterService {
if (msg.getEntityId().getEntityType().equals(EntityType.TENANT) if (msg.getEntityId().getEntityType().equals(EntityType.TENANT)
|| msg.getEntityId().getEntityType().equals(EntityType.TENANT_PROFILE) || msg.getEntityId().getEntityType().equals(EntityType.TENANT_PROFILE)
|| msg.getEntityId().getEntityType().equals(EntityType.DEVICE_PROFILE) || msg.getEntityId().getEntityType().equals(EntityType.DEVICE_PROFILE)
|| msg.getEntityId().getEntityType().equals(EntityType.API_USAGE_STATE)) { || msg.getEntityId().getEntityType().equals(EntityType.API_USAGE_STATE)
|| msg.getEntityId().getEntityType().equals(EntityType.EDGE)) {
TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer();
Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE); Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE);
for (String serviceId : tbCoreServices) { for (String serviceId : tbCoreServices) {

7
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -280,6 +280,13 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
} else if (toCoreNotification.getComponentLifecycleMsg() != null && !toCoreNotification.getComponentLifecycleMsg().isEmpty()) { } else if (toCoreNotification.getComponentLifecycleMsg() != null && !toCoreNotification.getComponentLifecycleMsg().isEmpty()) {
handleComponentLifecycleMsg(id, toCoreNotification.getComponentLifecycleMsg()); handleComponentLifecycleMsg(id, toCoreNotification.getComponentLifecycleMsg());
callback.onSuccess(); callback.onSuccess();
} else if (toCoreNotification.getEdgeEventUpdateMsg() != null && !toCoreNotification.getEdgeEventUpdateMsg().isEmpty()) {
Optional<TbActorMsg> actorMsg = encodingService.decode(toCoreNotification.getEdgeEventUpdateMsg().toByteArray());
if (actorMsg.isPresent()) {
log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get());
actorContext.tellWithHighPriority(actorMsg.get());
}
callback.onSuccess();
} }
if (statsEnabled) { if (statsEnabled) {
stats.log(toCoreNotification); stats.log(toCoreNotification);

3
application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java

@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
@ -71,4 +72,6 @@ public interface TbClusterService {
void onDeviceChange(Device device, TbQueueCallback callback); void onDeviceChange(Device device, TbQueueCallback callback);
void onDeviceDeleted(Device device, TbQueueCallback callback); void onDeviceDeleted(Device device, TbQueueCallback callback);
void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId);
} }

1
application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java

@ -52,7 +52,6 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi
private final TbClusterService clusterService; private final TbClusterService clusterService;
private final TbServiceInfoProvider serviceInfoProvider; private final TbServiceInfoProvider serviceInfoProvider;
private final ConcurrentMap<UUID, Consumer<FromDeviceRpcResponse>> toDeviceRpcRequests = new ConcurrentHashMap<>(); private final ConcurrentMap<UUID, Consumer<FromDeviceRpcResponse>> toDeviceRpcRequests = new ConcurrentHashMap<>();
private Optional<TbCoreDeviceRpcService> tbCoreRpcService; private Optional<TbCoreDeviceRpcService> tbCoreRpcService;

1
application/src/main/resources/thingsboard.yml

@ -589,6 +589,7 @@ edges:
max_read_records_count: "${EDGES_RPC_STORAGE_MAX_READ_RECORDS_COUNT:50}" max_read_records_count: "${EDGES_RPC_STORAGE_MAX_READ_RECORDS_COUNT:50}"
no_read_records_sleep: "${EDGES_RPC_NO_READ_RECORDS_SLEEP:1000}" no_read_records_sleep: "${EDGES_RPC_NO_READ_RECORDS_SLEEP:1000}"
sleep_between_batches: "${EDGES_RPC_SLEEP_BETWEEN_BATCHES:1000}" sleep_between_batches: "${EDGES_RPC_SLEEP_BETWEEN_BATCHES:1000}"
scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:4}"
edge_events_ttl: "${EDGES_EDGE_EVENTS_TTL:0}" edge_events_ttl: "${EDGES_EDGE_EVENTS_TTL:0}"
state: state:
persistToTelemetry: "${EDGES_PERSIST_STATE_TO_TELEMETRY:false}" persistToTelemetry: "${EDGES_PERSIST_STATE_TO_TELEMETRY:false}"

27
application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java

@ -98,6 +98,7 @@ import org.thingsboard.server.gen.edge.UserCredentialsUpdateMsg;
import org.thingsboard.server.gen.edge.WidgetTypeUpdateMsg; import org.thingsboard.server.gen.edge.WidgetTypeUpdateMsg;
import org.thingsboard.server.gen.edge.WidgetsBundleUpdateMsg; import org.thingsboard.server.gen.edge.WidgetsBundleUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.queue.TbClusterService;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
@ -122,6 +123,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
@Autowired @Autowired
private EdgeEventService edgeEventService; private EdgeEventService edgeEventService;
@Autowired
private TbClusterService clusterService;
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
loginSysAdmin(); loginSysAdmin();
@ -228,6 +232,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent); edgeEventService.saveAsync(edgeEvent);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage(); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -845,6 +850,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1); edgeEventService.saveAsync(edgeEvent1);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage(); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -883,6 +889,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_DELETED, device.getId().getId(), EdgeEventType.DEVICE, deleteAttributesEntityData); EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_DELETED, device.getId().getId(), EdgeEventType.DEVICE, deleteAttributesEntityData);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent); edgeEventService.saveAsync(edgeEvent);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage(); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -908,6 +915,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.POST_ATTRIBUTES, device.getId().getId(), EdgeEventType.DEVICE, postAttributesEntityData); EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.POST_ATTRIBUTES, device.getId().getId(), EdgeEventType.DEVICE, postAttributesEntityData);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent); edgeEventService.saveAsync(edgeEvent);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage(); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -932,6 +940,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1); edgeEventService.saveAsync(edgeEvent1);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage(); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -1158,6 +1167,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
edgeImitator.sendUplinkMsg(uplinkMsgBuilder2.build()); edgeImitator.sendUplinkMsg(uplinkMsgBuilder2.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
// Wait before device attributes saved to database before requesting them from controller
Thread.sleep(1000); Thread.sleep(1000);
Map<String, List<Map<String, String>>> timeseries = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/values/timeseries?keys=" + timeseriesKey, Map.class); Map<String, List<Map<String, String>>> timeseries = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/values/timeseries?keys=" + timeseriesKey, Map.class);
Assert.assertTrue(timeseries.containsKey(timeseriesKey)); Assert.assertTrue(timeseries.containsKey(timeseriesKey));
@ -1300,18 +1310,25 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void sendAttributesRequest() throws Exception { private void sendAttributesRequest() throws Exception {
Device device = findDeviceByName("Edge Device 1"); Device device = findDeviceByName("Edge Device 1");
sendAttributesRequest(device, DataConstants.SERVER_SCOPE, "{\"key1\":\"value1\"}", "key1", "value1");
sendAttributesRequest(device, DataConstants.SHARED_SCOPE, "{\"key2\":\"value2\"}", "key2", "value2");
}
String attributesDataStr = "{\"key1\":\"value1\"}"; private void sendAttributesRequest(Device device, String scope, String attributesDataStr, String expectedKey, String expectedValue) throws Exception {
JsonNode attributesData = mapper.readTree(attributesDataStr); JsonNode attributesData = mapper.readTree(attributesDataStr);
doPost("/api/plugins/telemetry/DEVICE/" + device.getId().getId().toString() + "/attributes/" + DataConstants.SERVER_SCOPE, doPost("/api/plugins/telemetry/DEVICE/" + device.getId().getId().toString() + "/attributes/" + scope,
attributesData); attributesData);
// Wait before device attributes saved to database before requesting them from edge
Thread.sleep(1000);
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
AttributesRequestMsg.Builder attributesRequestMsgBuilder = AttributesRequestMsg.newBuilder(); AttributesRequestMsg.Builder attributesRequestMsgBuilder = AttributesRequestMsg.newBuilder();
attributesRequestMsgBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits()); attributesRequestMsgBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits());
attributesRequestMsgBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits()); attributesRequestMsgBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits());
attributesRequestMsgBuilder.setEntityType(EntityType.DEVICE.name()); attributesRequestMsgBuilder.setEntityType(EntityType.DEVICE.name());
attributesRequestMsgBuilder.setScope(scope);
testAutoGeneratedCodeByProtobuf(attributesRequestMsgBuilder); testAutoGeneratedCodeByProtobuf(attributesRequestMsgBuilder);
uplinkMsgBuilder.addAttributesRequestMsg(attributesRequestMsgBuilder.build()); uplinkMsgBuilder.addAttributesRequestMsg(attributesRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
@ -1328,14 +1345,14 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB()); Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB());
Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB()); Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope()); Assert.assertEquals(scope, latestEntityDataMsg.getPostAttributeScope());
Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg()); Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg());
TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg(); TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg();
Assert.assertEquals(1, attributesUpdatedMsg.getKvCount()); Assert.assertEquals(1, attributesUpdatedMsg.getKvCount());
TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0); TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0);
Assert.assertEquals("key1", keyValueProto.getKey()); Assert.assertEquals(expectedKey, keyValueProto.getKey());
Assert.assertEquals("value1", keyValueProto.getStringV()); Assert.assertEquals(expectedValue, keyValueProto.getStringV());
} }
private void sendDeleteDeviceOnEdge() throws Exception { private void sendDeleteDeviceOnEdge() throws Exception {

22
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -56,6 +56,8 @@ import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@Slf4j @Slf4j
public class EdgeImitator { public class EdgeImitator {
@ -65,6 +67,8 @@ public class EdgeImitator {
private EdgeRpcClient edgeRpcClient; private EdgeRpcClient edgeRpcClient;
private final Lock lock = new ReentrantLock();
private CountDownLatch messagesLatch; private CountDownLatch messagesLatch;
private CountDownLatch responsesLatch; private CountDownLatch responsesLatch;
private List<Class<? extends AbstractMessage>> ignoredTypes; private List<Class<? extends AbstractMessage>> ignoredTypes;
@ -74,7 +78,7 @@ public class EdgeImitator {
@Getter @Getter
private UserId userId; private UserId userId;
@Getter @Getter
private List<AbstractMessage> downlinkMsgs; private final List<AbstractMessage> downlinkMsgs;
public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException { public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException {
edgeRpcClient = new EdgeGrpcClient(); edgeRpcClient = new EdgeGrpcClient();
@ -241,7 +245,12 @@ public class EdgeImitator {
private ListenableFuture<Void> saveDownlinkMsg(AbstractMessage message) { private ListenableFuture<Void> saveDownlinkMsg(AbstractMessage message) {
if (!ignoredTypes.contains(message.getClass())) { if (!ignoredTypes.contains(message.getClass())) {
downlinkMsgs.add(message); try {
lock.lock();
downlinkMsgs.add(message);
} finally {
lock.unlock();
}
messagesLatch.countDown(); messagesLatch.countDown();
} }
return Futures.immediateFuture(null); return Futures.immediateFuture(null);
@ -262,7 +271,14 @@ public class EdgeImitator {
} }
public <T> Optional<T> findMessageByType(Class<T> tClass) { public <T> Optional<T> findMessageByType(Class<T> tClass) {
return (Optional<T>) downlinkMsgs.stream().filter(downlinkMsg -> downlinkMsg.getClass().isAssignableFrom(tClass)).findAny(); Optional<T> result;
try {
lock.lock();
result = (Optional<T>) downlinkMsgs.stream().filter(downlinkMsg -> downlinkMsg.getClass().isAssignableFrom(tClass)).findAny();
} finally {
lock.unlock();
}
return result;
} }
public AbstractMessage getLatestMessage() { public AbstractMessage getLatestMessage() {

3
common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java

@ -15,8 +15,10 @@
*/ */
package org.thingsboard.server.common.data; package org.thingsboard.server.common.data;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.edge.EdgeEventType;
@Slf4j
public final class EdgeUtils { public final class EdgeUtils {
private EdgeUtils() { private EdgeUtils() {
@ -49,6 +51,7 @@ public final class EdgeUtils {
case WIDGET_TYPE: case WIDGET_TYPE:
return EdgeEventType.WIDGET_TYPE; return EdgeEventType.WIDGET_TYPE;
default: default:
log.warn("Unsupported entity type [{}]", entityType);
return null; return null;
} }
} }

17
common/edge-api/src/main/proto/edge.proto

@ -79,12 +79,16 @@ message EdgeConfiguration {
int64 edgeIdLSB = 2; int64 edgeIdLSB = 2;
int64 tenantIdMSB = 3; int64 tenantIdMSB = 3;
int64 tenantIdLSB = 4; int64 tenantIdLSB = 4;
string name = 5; int64 customerIdMSB = 5;
string routingKey = 6; int64 customerIdLSB = 6;
string type = 7; string name = 7;
string edgeLicenseKey = 8; string type = 8;
string cloudEndpoint = 9; string routingKey = 9;
string cloudType = 10; string secret = 10;
string edgeLicenseKey = 11;
string cloudEndpoint = 12;
string configuration = 13;
string cloudType = 14;
} }
enum UpdateMsgType { enum UpdateMsgType {
@ -315,6 +319,7 @@ message AttributesRequestMsg {
int64 entityIdMSB = 1; int64 entityIdMSB = 1;
int64 entityIdLSB = 2; int64 entityIdLSB = 2;
string entityType = 3; string entityType = 3;
string scope = 4;
} }
message RelationRequestMsg { message RelationRequestMsg {

7
common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java

@ -105,6 +105,11 @@ public enum MsgType {
/** /**
* Message that is sent by TransportRuleEngineService to Device Actor. Represents messages from the device itself. * Message that is sent by TransportRuleEngineService to Device Actor. Represents messages from the device itself.
*/ */
TRANSPORT_TO_DEVICE_ACTOR_MSG; TRANSPORT_TO_DEVICE_ACTOR_MSG,
/**
* Message that is sent on Edge Event to Edge Session
*/
EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG;
} }

42
common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeEventUpdateMsg.java

@ -0,0 +1,42 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.msg.edge;
import lombok.Getter;
import lombok.ToString;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.aware.TenantAwareMsg;
import org.thingsboard.server.common.msg.cluster.ToAllNodesMsg;
@ToString
public class EdgeEventUpdateMsg implements TenantAwareMsg, ToAllNodesMsg {
@Getter
private final TenantId tenantId;
@Getter
private final EdgeId edgeId;
public EdgeEventUpdateMsg(TenantId tenantId, EdgeId edgeId) {
this.tenantId = tenantId;
this.edgeId = edgeId;
}
@Override
public MsgType getMsgType() {
return MsgType.EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG;
}
}

1
common/queue/src/main/proto/queue.proto

@ -520,6 +520,7 @@ message ToCoreNotificationMsg {
LocalSubscriptionServiceMsgProto toLocalSubscriptionServiceMsg = 1; LocalSubscriptionServiceMsgProto toLocalSubscriptionServiceMsg = 1;
FromDeviceRPCResponseProto fromDeviceRpcResponse = 2; FromDeviceRPCResponseProto fromDeviceRpcResponse = 2;
bytes componentLifecycleMsg = 3; bytes componentLifecycleMsg = 3;
bytes edgeEventUpdateMsg = 4;
} }
/* Messages that are handled by ThingsBoard RuleEngine Service */ /* Messages that are handled by ThingsBoard RuleEngine Service */

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

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.edge; package org.thingsboard.server.dao.edge;
import com.google.common.base.Function; import com.google.common.base.Function;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
@ -98,6 +99,8 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId "; public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId ";
public static final String INCORRECT_EDGE_ID = "Incorrect edgeId "; public static final String INCORRECT_EDGE_ID = "Incorrect edgeId ";
private static final int DEFAULT_LIMIT = 100;
private RestTemplate restTemplate; private RestTemplate restTemplate;
private static final String EDGE_LICENSE_SERVER_ENDPOINT = "https://license.thingsboard.io"; private static final String EDGE_LICENSE_SERVER_ENDPOINT = "https://license.thingsboard.io";
@ -353,13 +356,20 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
public void assignDefaultRuleChainsToEdge(TenantId tenantId, EdgeId edgeId) { public void assignDefaultRuleChainsToEdge(TenantId tenantId, EdgeId edgeId) {
log.trace("Executing assignDefaultRuleChainsToEdge, tenantId [{}], edgeId [{}]", tenantId, edgeId); log.trace("Executing assignDefaultRuleChainsToEdge, tenantId [{}], edgeId [{}]", tenantId, edgeId);
ListenableFuture<List<RuleChain>> future = ruleChainService.findDefaultEdgeRuleChainsByTenantId(tenantId); ListenableFuture<List<RuleChain>> future = ruleChainService.findDefaultEdgeRuleChainsByTenantId(tenantId);
Futures.transform(future, ruleChains -> { Futures.addCallback(future, new FutureCallback<List<RuleChain>>() {
if (ruleChains != null && !ruleChains.isEmpty()) { @Override
for (RuleChain ruleChain : ruleChains) { public void onSuccess(List<RuleChain> ruleChains) {
ruleChainService.assignRuleChainToEdge(tenantId, ruleChain.getId(), edgeId); if (ruleChains != null && !ruleChains.isEmpty()) {
for (RuleChain ruleChain : ruleChains) {
ruleChainService.assignRuleChainToEdge(tenantId, ruleChain.getId(), edgeId);
}
} }
} }
return null;
@Override
public void onFailure(Throwable t) {
log.warn("[{}] can't find default edge rule chains [{}]", tenantId.getId(), edgeId.getId(), t);
}
}, MoreExecutors.directExecutor()); }, MoreExecutors.directExecutor());
} }
@ -462,9 +472,26 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
@Override @Override
public ListenableFuture<List<EdgeId>> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) { public ListenableFuture<List<EdgeId>> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) {
log.trace("[{}] Executing findRelatedEdgeIdsByEntityId [{}]", tenantId, entityId); log.trace("[{}] Executing findRelatedEdgeIdsByEntityId [{}]", tenantId, entityId);
if (EntityType.TENANT.equals(entityId.getEntityType())) { if (EntityType.TENANT.equals(entityId.getEntityType()) || EntityType.CUSTOMER.equals(entityId.getEntityType())) {
PageData<Edge> edgesByTenantId = findEdgesByTenantId(tenantId, new PageLink(Integer.MAX_VALUE)); List<EdgeId> result = new ArrayList<>();
return Futures.immediateFuture(edgesByTenantId.getData().stream().map(IdBased::getId).collect(Collectors.toList())); PageLink pageLink = new PageLink(DEFAULT_LIMIT);
PageData<Edge> pageData;
do {
if (EntityType.TENANT.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);
} else { } else {
switch (entityId.getEntityType()) { switch (entityId.getEntityType()) {
case DEVICE: case DEVICE:
@ -489,13 +516,23 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
if (userById == null) { if (userById == null) {
return Futures.immediateFuture(Collections.emptyList()); return Futures.immediateFuture(Collections.emptyList());
} }
PageData<Edge> edges; List<Edge> result = new ArrayList<>();
if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) { PageLink pageLink = new PageLink(DEFAULT_LIMIT);
edges = findEdgesByTenantId(tenantId, new PageLink(Integer.MAX_VALUE)); PageData<Edge> pageData;
} else { do {
edges = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), new PageLink(Integer.MAX_VALUE)); if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) {
} pageData = findEdgesByTenantId(tenantId, pageLink);
return convertToEdgeIds(Futures.immediateFuture(edges.getData())); } else {
pageData = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), 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: default:
return Futures.immediateFuture(Collections.emptyList()); return Futures.immediateFuture(Collections.emptyList());
} }

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

@ -77,6 +77,7 @@ public abstract class AbstractEntityService {
List<EntityView> entityViews = entityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId).get(); List<EntityView> entityViews = entityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId).get();
if (entityViews != null && !entityViews.isEmpty()) { if (entityViews != null && !entityViews.isEmpty()) {
EntityView entityView = entityViews.get(0); EntityView entityView = entityViews.get(0);
// TODO: voba - refactor this blocking operation in 3.3+
Boolean relationExists = relationService.checkRelation(tenantId,edgeId, entityView.getId(), Boolean relationExists = relationService.checkRelation(tenantId,edgeId, entityView.getId(),
EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE).get(); EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE).get();
if (relationExists) { if (relationExists) {
@ -84,7 +85,7 @@ public abstract class AbstractEntityService {
} }
} }
} catch (Exception e) { } catch (Exception e) {
log.error("Exception while finding entity views for entityId [{}]", entityId, e); log.error("[{}] Exception while finding entity views for entityId [{}]", tenantId, entityId, e);
throw new RuntimeException("Exception while finding entity views for entityId [" + entityId + "]", e); throw new RuntimeException("Exception while finding entity views for entityId [" + entityId + "]", e);
} }
} }

16
dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java

@ -59,7 +59,7 @@ public class BaseEventService implements EventService {
public Optional<Event> saveIfNotExists(Event event) { public Optional<Event> saveIfNotExists(Event event) {
eventValidator.validate(event, Event::getTenantId); eventValidator.validate(event, Event::getTenantId);
if (StringUtils.isEmpty(event.getUid())) { if (StringUtils.isEmpty(event.getUid())) {
throw new DataValidationException("Event uid should be specified!"); throw new DataValidationException("Event uid should be specified!.");
} }
checkAndTruncateDebugEvent(event); checkAndTruncateDebugEvent(event);
return eventDao.saveIfNotExists(event); return eventDao.saveIfNotExists(event);
@ -79,16 +79,16 @@ public class BaseEventService implements EventService {
@Override @Override
public Optional<Event> findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) { public Optional<Event> findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) {
if (tenantId == null) { if (tenantId == null) {
throw new DataValidationException("Tenant id should be specified!"); throw new DataValidationException("Tenant id should be specified!.");
} }
if (entityId == null) { if (entityId == null) {
throw new DataValidationException("Entity id should be specified!"); throw new DataValidationException("Entity id should be specified!.");
} }
if (StringUtils.isEmpty(eventType)) { if (StringUtils.isEmpty(eventType)) {
throw new DataValidationException("Event type should be specified!"); throw new DataValidationException("Event type should be specified!.");
} }
if (StringUtils.isEmpty(eventUid)) { if (StringUtils.isEmpty(eventUid)) {
throw new DataValidationException("Event uid should be specified!"); throw new DataValidationException("Event uid should be specified!.");
} }
Event event = eventDao.findEvent(tenantId.getId(), entityId, eventType, eventUid); Event event = eventDao.findEvent(tenantId.getId(), entityId, eventType, eventUid);
return event != null ? Optional.of(event) : Optional.empty(); return event != null ? Optional.of(event) : Optional.empty();
@ -129,13 +129,13 @@ public class BaseEventService implements EventService {
@Override @Override
protected void validateDataImpl(TenantId tenantId, Event event) { protected void validateDataImpl(TenantId tenantId, Event event) {
if (event.getEntityId() == null) { if (event.getEntityId() == null) {
throw new DataValidationException("Entity id should be specified!"); throw new DataValidationException("Entity id should be specified!.");
} }
if (StringUtils.isEmpty(event.getType())) { if (StringUtils.isEmpty(event.getType())) {
throw new DataValidationException("Event type should be specified!"); throw new DataValidationException("Event type should be specified!.");
} }
if (event.getBody() == null) { if (event.getBody() == null) {
throw new DataValidationException("Event body should be specified!"); throw new DataValidationException("Event body should be specified!.");
} }
} }
}; };

3
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -27,6 +27,7 @@ import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -153,6 +154,8 @@ public interface TbContext {
// TODO: Does this changes the message? // TODO: Does this changes the message?
TbMsg alarmActionMsg(Alarm alarm, RuleNodeId ruleNodeId, String action); TbMsg alarmActionMsg(Alarm alarm, RuleNodeId ruleNodeId, String action);
void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId);
/* /*
* *
* METHODS TO PROCESS THE MESSAGES * METHODS TO PROCESS THE MESSAGES

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

@ -57,7 +57,7 @@ import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
name = "push to edge", name = "push to edge",
configClazz = EmptyNodeConfiguration.class, configClazz = EmptyNodeConfiguration.class,
nodeDescription = "Pushes messages to edge", nodeDescription = "Pushes messages to edge",
nodeDetails = "Pushes messages to edge, if Message Originator assigned to particular edge or is EDGE entity. This node is used only on Cloud instances to push messages from Cloud to Edge. Supports only DEVICE, ENTITY_VIEW, ASSET and EDGE Message Originator(s).", nodeDetails = "Pushes messages to edge, if Message Originator assigned to particular edge or is EDGE entity. This node is used only on Cloud instances to push messages from Cloud to Edge. Supports only DEVICE, ENTITY_VIEW, ASSET, ENTITY_VIEW, DASHBOARD, TENANT, CUSTOMER and EDGE Message Originator(s).",
uiResources = {"static/rulenode/rulenode-core-config.js", "static/rulenode/rulenode-core-config.css"}, uiResources = {"static/rulenode/rulenode-core-config.js", "static/rulenode/rulenode-core-config.css"},
configDirective = "tbNodeEmptyConfig", configDirective = "tbNodeEmptyConfig",
icon = "cloud_download", icon = "cloud_download",
@ -97,47 +97,75 @@ public class TbMsgPushToEdgeNode implements TbNode {
} }
private void processMsg(TbContext ctx, TbMsg msg) { private void processMsg(TbContext ctx, TbMsg msg) {
ListenableFuture<List<EdgeId>> getEdgeIdsFuture = ctx.getEdgeService().findRelatedEdgeIdsByEntityId(ctx.getTenantId(), msg.getOriginator()); if (EntityType.EDGE.equals(msg.getOriginator().getEntityType())) {
Futures.addCallback(getEdgeIdsFuture, new FutureCallback<List<EdgeId>>() { try {
@Override EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx);
public void onSuccess(@Nullable List<EdgeId> edgeIds) { if (edgeEvent != null) {
if (edgeIds != null && !edgeIds.isEmpty()) { EdgeId edgeId = new EdgeId(msg.getOriginator().getId());
for (EdgeId edgeId : edgeIds) { edgeEvent.setEdgeId(edgeId);
try { ListenableFuture<EdgeEvent> saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent);
EdgeEvent edgeEvent = buildEdgeEvent(msg, ctx); Futures.addCallback(saveFuture, new FutureCallback<EdgeEvent>() {
if (edgeEvent == null) { @Override
log.debug("Edge event type is null. Entity Type {}", msg.getOriginator().getEntityType()); public void onSuccess(@Nullable EdgeEvent event) {
ctx.tellFailure(msg, new RuntimeException("Edge event type is null. Entity Type '" + msg.getOriginator().getEntityType() + "'")); ctx.tellNext(msg, SUCCESS);
} else { ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId);
edgeEvent.setEdgeId(edgeId); }
ListenableFuture<EdgeEvent> saveFuture = ctx.getEdgeEventService().saveAsync(edgeEvent);
Futures.addCallback(saveFuture, new FutureCallback<EdgeEvent>() { @Override
@Override public void onFailure(Throwable th) {
public void onSuccess(@Nullable EdgeEvent event) { log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th);
ctx.tellNext(msg, SUCCESS); ctx.tellFailure(msg, th);
} }
}, ctx.getDbCallbackExecutor());
@Override }
public void onFailure(Throwable th) { } catch (JsonProcessingException e) {
log.error("Could not save edge event", th); log.error("Failed to build edge event", e);
ctx.tellFailure(msg, th); ctx.tellFailure(msg, e);
} }
}, ctx.getDbCallbackExecutor()); } 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);
} }
} catch (JsonProcessingException e) {
log.error("Failed to build edge event", e);
ctx.tellFailure(msg, e);
} }
} }
} }
}
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
ctx.tellFailure(msg, t); ctx.tellFailure(msg, t);
} }
}, ctx.getDbCallbackExecutor()); }, ctx.getDbCallbackExecutor());
}
} }
private EdgeEvent buildEdgeEvent(TbMsg msg, TbContext ctx) throws JsonProcessingException { private EdgeEvent buildEdgeEvent(TbMsg msg, TbContext ctx) throws JsonProcessingException {
@ -220,6 +248,7 @@ public class TbMsgPushToEdgeNode implements TbNode {
case DASHBOARD: case DASHBOARD:
case TENANT: case TENANT:
case CUSTOMER: case CUSTOMER:
case EDGE:
return true; return true;
default: default:
return false; return false;

14
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -307,10 +307,11 @@
"filter-type-device-type-description": "Devices of type '{{deviceType}}'", "filter-type-device-type-description": "Devices of type '{{deviceType}}'",
"filter-type-device-type-and-name-description": "Devices of type '{{deviceType}}' and with name starting with '{{prefix}}'", "filter-type-device-type-and-name-description": "Devices of type '{{deviceType}}' and with name starting with '{{prefix}}'",
"filter-type-entity-view-type": "Entity View type", "filter-type-entity-view-type": "Entity View type",
"filter-type-entity-view-type-description": "Entity Views of type '{{entityView}}'", "filter-type-entity-view-type-description": "Entity Views of type '{{entityViewType}}'",
"filter-type-entity-view-type-and-name-description": "Entity Views of type '{{entityView}}' and with name starting with '{{prefix}}'", "filter-type-entity-view-type-and-name-description": "Entity Views of type '{{entityViewType}}' and with name starting with '{{prefix}}'",
"filter-type-edge-type": "Edge type", "filter-type-edge-type": "Edge type",
"filter-type-edge-type-description": "Edges of type '{{edgeType}}'", "filter-type-edge-type-description": "Edges of type '{{edgeType}}'",
"filter-type-edge-type-and-name-description": "Edges of type '{{edgeType}}' and with name starting with '{{prefix}}'",
"filter-type-relations-query": "Relations query", "filter-type-relations-query": "Relations query",
"filter-type-relations-query-description": "{{entities}} that have {{relationType}} relation {{direction}} {{rootEntity}}", "filter-type-relations-query-description": "{{entities}} that have {{relationType}} relation {{direction}} {{rootEntity}}",
"filter-type-asset-search-query": "Asset search query", "filter-type-asset-search-query": "Asset search query",
@ -1131,6 +1132,7 @@
"delete-edges-action-title": "Delete { count, plural, 1 {1 edge} other {# edges} }", "delete-edges-action-title": "Delete { count, plural, 1 {1 edge} other {# edges} }",
"delete-edges-text": "Be careful, after the confirmation all selected edges will be removed and all related data will become unrecoverable.", "delete-edges-text": "Be careful, after the confirmation all selected edges will be removed and all related data will become unrecoverable.",
"name": "Name", "name": "Name",
"name-starts-with": "Edge name starts with",
"name-required": "Name is required.", "name-required": "Name is required.",
"edge-license-key": "Edge License Key", "edge-license-key": "Edge License Key",
"edge-license-key-required": "Edge License Key is required.", "edge-license-key-required": "Edge License Key is required.",
@ -1201,7 +1203,8 @@
"selected-edges": "{ count, plural, 1 {1 edge} other {# edges} } selected", "selected-edges": "{ count, plural, 1 {1 edge} other {# edges} } selected",
"enter-edge-type": "Enter entity view type", "enter-edge-type": "Enter entity view type",
"any-edge": "Any edge", "any-edge": "Any edge",
"no-edge-types-matching": "No edge types matching '{{entitySubtype}}' were found." "no-edge-types-matching": "No edge types matching '{{entitySubtype}}' were found.",
"unassign-edges-action-title": "Unassign { count, plural, 1 {1 edge} other {# edges} } from customer"
}, },
"error": { "error": {
"unable-to-connect": "Unable to connect to the server! Please check your internet connection.", "unable-to-connect": "Unable to connect to the server! Please check your internet connection.",
@ -1993,6 +1996,7 @@
"rulechain": { "rulechain": {
"rulechain": "Rule chain", "rulechain": "Rule chain",
"rulechains": "Rule chains", "rulechains": "Rule chains",
"default-root": "Default root",
"root": "Root", "root": "Root",
"delete": "Delete rule chain", "delete": "Delete rule chain",
"name": "Name", "name": "Name",
@ -2041,10 +2045,10 @@
"set-default-root-edge-rulechain-title": "Are you sure you want to make the rule chain '{{ruleChainName}}' default edge root?", "set-default-root-edge-rulechain-title": "Are you sure you want to make the rule chain '{{ruleChainName}}' default edge root?",
"set-default-root-edge-rulechain-text": "After the confirmation the rule chain will become default edge root and will handle all incoming transport messages.", "set-default-root-edge-rulechain-text": "After the confirmation the rule chain will become default edge root and will handle all incoming transport messages.",
"invalid-rulechain-type-error": "Unable to import rule chain: Invalid rule chain type. Expected type is {{expectedRuleChainType}}.", "invalid-rulechain-type-error": "Unable to import rule chain: Invalid rule chain type. Expected type is {{expectedRuleChainType}}.",
"set-default-edge": "Make edge rule chain default", "set-default-edge": "Make rule chain default",
"set-default-edge-title": "Are you sure you want to make the edge rule chain '{{ruleChainName}}' default?", "set-default-edge-title": "Are you sure you want to make the edge rule chain '{{ruleChainName}}' default?",
"set-default-edge-text": "After the confirmation the edge rule chain will be added to default list and assigned to newly created edge(s).", "set-default-edge-text": "After the confirmation the edge rule chain will be added to default list and assigned to newly created edge(s).",
"remove-default-edge": "Remove edge rule chain from defaults", "remove-default-edge": "Remove rule chain from defaults",
"remove-default-edge-title": "Are you sure you want to remove the edge rule chain '{{ruleChainName}}' from default list?", "remove-default-edge-title": "Are you sure you want to remove the edge rule chain '{{ruleChainName}}' from default list?",
"remove-default-edge-text": "After the confirmation the edge rule chain will not be assigned for a newly created edges.", "remove-default-edge-text": "After the confirmation the edge rule chain will not be assigned for a newly created edges.",
"unassign-rulechain-title": "Are you sure you want to unassign the rulechain '{{ruleChainName}}'?", "unassign-rulechain-title": "Are you sure you want to unassign the rulechain '{{ruleChainName}}'?",

2
ui/src/app/edge/edge.controller.js

@ -283,7 +283,7 @@ export function EdgeController($rootScope, userService, edgeService, customerSer
details: function() { details: function() {
return $translate.instant('edge.manage-edge-rulechains'); return $translate.instant('edge.manage-edge-rulechains');
}, },
icon: "settings_ethernet" icon: "code"
} }
); );

2
ui/src/app/edge/edge.routes.js

@ -49,7 +49,7 @@ export default function EdgeRoutes($stateProvider, types) {
pageTitle: 'edge.edges' pageTitle: 'edge.edges'
}, },
ncyBreadcrumb: { ncyBreadcrumb: {
label: '{"icon": "transform", "label": "edge.edges"}' label: '{"icon": "router", "label": "edge.edges"}'
} }
}).state('home.edges.entityViews', { }).state('home.edges.entityViews', {
url: '/:edgeId/entityViews', url: '/:edgeId/entityViews',

Loading…
Cancel
Save