Browse Source

Refactored Edge session message poll mechanism. Removed tech dept - Integer.MAX_INTEGER. Code refactoring

pull/3693/head
Volodymyr Babak 6 years ago
parent
commit
193cccf123
  1. 18
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  2. 20
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  3. 6
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  4. 9
      application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
  5. 1
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  6. 7
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  7. 113
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  8. 38
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  9. 5
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  10. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeRpcService.java
  11. 121
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java
  12. 17
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseProcessor.java
  13. 36
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  14. 18
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  15. 3
      application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java
  16. 1
      application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java
  17. 27
      application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java
  18. 22
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  19. 1
      common/edge-api/src/main/proto/edge.proto
  20. 7
      common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java
  21. 42
      common/message/src/main/java/org/thingsboard/server/common/msg/edge/EdgeEventUpdateMsg.java
  22. 1
      common/queue/src/main/proto/queue.proto
  23. 43
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  24. 3
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  25. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

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

@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg;
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.queue.QueueToRuleEngineMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
@ -93,6 +94,9 @@ public class AppActor extends ContextAwareActor {
case SERVER_RPC_RESPONSE_TO_DEVICE_ACTOR_MSG:
onToDeviceActorMsg((TenantAwareMsg) msg, true);
break;
case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG:
onToTenantActorMsg((EdgeEventUpdateMsg) msg);
break;
default:
return false;
}
@ -186,6 +190,20 @@ public class AppActor extends ContextAwareActor {
() -> new TenantActor.ActorCreator(systemContext, tenantId));
}
private void onToTenantActorMsg(EdgeEventUpdateMsg msg) {
TbActorRef target = null;
if (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 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) {
EntityRelation relationToEdge = result.get(0);
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());
} else {
log.trace("[{}][{}] edge relation is empty {}", tenantId, deviceId, relationToEdge);
}
} else {
log.trace("[{}][{}] device doesn't have any related edge", tenantId, deviceId);
}
return null;
}
@ -165,6 +170,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
boolean sent;
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());
sent = true;
} else {
@ -516,6 +522,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
}
void processEdgeUpdate(DeviceEdgeUpdateMsg msg) {
log.trace("[{}] Processing edge update {}", deviceId, msg);
this.edgeId = msg.getEdgeId();
}
@ -568,7 +575,18 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
edgeEvent.setBody(body);
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) {

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

@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId;
@ -269,6 +270,11 @@ class DefaultTbContext implements TbContext {
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) {
try {
return TbMsg.newMsg(action, id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)));

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

@ -45,6 +45,7 @@ import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.aware.DeviceAwareMsg;
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.queue.PartitionChangeMsg;
import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
@ -157,6 +158,9 @@ public class TenantActor extends RuleChainManagerActor {
case RULE_CHAIN_TO_RULE_CHAIN_MSG:
onRuleChainMsg((RuleChainAwareMsg) msg);
break;
case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG:
onToEdgeSessionMsg((EdgeEventUpdateMsg) msg);
break;
default:
return false;
}
@ -242,6 +246,11 @@ public class TenantActor extends RuleChainManagerActor {
() -> 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 {
private final TenantId tenantId;

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

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

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

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

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

@ -55,6 +55,7 @@ import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.queue.TbClusterService;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
@ -74,6 +75,8 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
private static final ObjectMapper mapper = new ObjectMapper();
private static final int DEFAULT_LIMIT = 100;
@Autowired
private EdgeService edgeService;
@ -89,6 +92,9 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
@Autowired
private EdgeEventService edgeEventService;
@Autowired
private TbClusterService clusterService;
@Autowired
private DbCallbackExecutorService dbCallbackExecutorService;
@ -137,7 +143,19 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
edgeEvent.setEntityId(entityId.getId());
}
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
@ -195,13 +213,20 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
public void onSuccess(@Nullable Edge edge) {
if (edge != null && !customerId.isNullUid()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, customerId, null);
TextPageData<User> pageData = userService.findCustomerUsers(tenantId, customerId, new TextPageLink(Integer.MAX_VALUE));
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);
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<User> pageData;
do {
pageData = userService.findCustomerUsers(tenantId, customerId, pageLink);
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 = pageData.getNextPageLink();
}
}
}
} while (pageData != null && pageData.hasNext());
}
}
@ -242,12 +267,19 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
case ADDED:
case UPDATED:
case DELETED:
TextPageData<Edge> edgesByTenantId = edgeService.findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE));
if (edgesByTenantId != null && edgesByTenantId.getData() != null && !edgesByTenantId.getData().isEmpty()) {
for (Edge edge : edgesByTenantId.getData()) {
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null);
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<Edge> pageData;
do {
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 = pageData.getNextPageLink();
}
}
}
} while (pageData != null && pageData.hasNext());
break;
}
}
@ -256,21 +288,28 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction());
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType());
EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB()));
TextPageData<Edge> edgesByTenantId = edgeService.findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE));
if (edgesByTenantId != null && edgesByTenantId.getData() != null && !edgesByTenantId.getData().isEmpty()) {
for (Edge edge : edgesByTenantId.getData()) {
switch (actionType) {
case UPDATED:
if (!edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(entityId)) {
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<Edge> pageData;
do {
pageData = edgeService.findEdgesByTenantId(tenantId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
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);
}
break;
case DELETED:
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null);
break;
break;
}
}
if (pageData.hasNext()) {
pageLink = pageData.getNextPageLink();
}
}
}
} while (pageData != null && pageData.hasNext());
}
private void processEntity(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) {
@ -337,26 +376,33 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
}, dbCallbackExecutorService);
break;
case DELETED:
TextPageData<Edge> edgesByTenantId = edgeService.findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE));
if (edgesByTenantId != null && edgesByTenantId.getData() != null && !edgesByTenantId.getData().isEmpty()) {
for (Edge edge : edgesByTenantId.getData()) {
saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, null);
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<Edge> pageData;
do {
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 = pageData.getNextPageLink();
}
}
}
} while (pageData != null && pageData.hasNext());
break;
case ASSIGNED_TO_EDGE:
case UNASSIGNED_FROM_EDGE:
EdgeId edgeId = new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB()));
saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null);
if (type.equals(EdgeEventType.RULE_CHAIN)) {
updateDependentRuleChains(tenantId, new RuleChainId(entityId.getId()), edgeId);
updateDependentRuleChains(tenantId, new RuleChainId(entityId.getId()), edgeId, new TimePageLink(DEFAULT_LIMIT));
}
break;
}
}
private void updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) {
ListenableFuture<TimePageData<RuleChain>> future = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, new TimePageLink(Integer.MAX_VALUE));
private void updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId, TimePageLink pageLink) {
ListenableFuture<TimePageData<RuleChain>> future = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<RuleChain>>() {
@Override
public void onSuccess(@Nullable TimePageData<RuleChain> pageData) {
@ -379,6 +425,9 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
}
}
}
if (pageData.hasNext()) {
updateDependentRuleChains(tenantId, processingRuleChainId, edgeId, pageData.getNextPageLink());
}
}
}

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

@ -49,6 +49,7 @@ import java.io.IOException;
import java.util.Collections;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
@ -59,7 +60,8 @@ import java.util.concurrent.TimeUnit;
@TbCoreComponent
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 static final ObjectMapper mapper = new ObjectMapper();
@Value("${edges.rpc.port}")
@ -147,12 +149,23 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
log.debug("Closing and removing session for edge [{}]", edgeId);
session.close();
sessions.remove(edgeId);
sessionNewEvents.remove(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) {
log.debug("[{}] onEdgeConnect [{}]", edgeId, edgeGrpcSession.getSessionId());
sessions.put(edgeId, edgeGrpcSession);
sessionNewEvents.put(edgeId, false);
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true);
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis());
}
@ -171,15 +184,23 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
while (!Thread.interrupted()) {
try {
if (sessions.size() > 0) {
for (EdgeGrpcSession session : sessions.values()) {
session.processHandleMessages();
for (Map.Entry<EdgeId, EdgeGrpcSession> entry : sessions.entrySet()) {
EdgeId edgeId = entry.getKey();
EdgeGrpcSession session = entry.getValue();
if (sessionNewEvents.get(edgeId)) {
log.trace("[{}] set session new events flag to false", edgeId.getId());
sessionNewEvents.put(edgeId, false);
// TODO: voba - at the moment all edge events are processed in a single thread. Maybe this should be updated?
session.processHandleMessages();
}
}
} else {
log.trace("No sessions available, sleep for the next run");
try {
Thread.sleep(1000);
} catch (InterruptedException ignore) {
}
log.trace("No sessions available");
}
log.trace("Sleep for the next run");
try {
Thread.sleep(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval());
} catch (InterruptedException ignore) {
}
} catch (Exception e) {
log.warn("Failed to process messages handling!", e);
@ -195,6 +216,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private void onEdgeDisconnect(EdgeId edgeId) {
log.debug("[{}] onEdgeDisconnect", edgeId);
sessions.remove(edgeId);
sessionNewEvents.remove(edgeId);
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false);
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis());
}

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

@ -304,11 +304,6 @@ public final class EdgeGrpcSession implements Closeable {
Long newStartTs = UUIDs.unixTimestamp(ifOffset);
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);
}

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 deleteEdge(EdgeId edgeId);
void onEdgeEvent(EdgeId edgeId);
}

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

@ -82,6 +82,7 @@ import org.thingsboard.server.gen.edge.RelationRequestMsg;
import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg;
import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.queue.TbClusterService;
import java.io.File;
import java.nio.charset.StandardCharsets;
@ -99,6 +100,8 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private static final ObjectMapper mapper = new ObjectMapper();
private static final int DEFAULT_LIMIT = 100;
@Autowired
private EdgeEventService edgeEventService;
@ -138,28 +141,31 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
@Autowired
private DbCallbackExecutorService dbCallbackExecutorService;
@Autowired
private TbClusterService tbClusterService;
@Override
public void sync(Edge edge) {
log.trace("[{}] staring sync process for edge [{}]", edge.getTenantId(), edge.getName());
try {
syncWidgetsBundleAndWidgetTypes(edge);
syncAdminSettings(edge);
syncRuleChains(edge);
syncRuleChains(edge, new TimePageLink(DEFAULT_LIMIT));
syncUsers(edge);
syncDevices(edge);
syncAssets(edge);
syncEntityViews(edge);
syncDashboards(edge);
syncDevices(edge, new TimePageLink(DEFAULT_LIMIT));
syncAssets(edge, new TimePageLink(DEFAULT_LIMIT));
syncEntityViews(edge, new TimePageLink(DEFAULT_LIMIT));
syncDashboards(edge, new TimePageLink(DEFAULT_LIMIT));
} catch (Exception e) {
log.error("Exception during sync process", e);
}
}
private void syncRuleChains(Edge edge) {
log.trace("[{}] syncRuleChains [{}]", edge.getTenantId(), edge.getName());
private void syncRuleChains(Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncRuleChains [{}] [{}]", edge.getTenantId(), edge.getName(), pageLink);
try {
ListenableFuture<TimePageData<RuleChain>> future =
ruleChainService.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE));
ruleChainService.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<RuleChain>>() {
@Override
public void onSuccess(@Nullable TimePageData<RuleChain> pageData) {
@ -168,6 +174,9 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
for (RuleChain ruleChain : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.RULE_CHAIN, EdgeEventActionType.ADDED, ruleChain.getId(), null);
}
if (pageData.hasNext()) {
syncRuleChains(edge, pageData.getNextPageLink());
}
}
}
@ -181,11 +190,11 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncDevices(Edge edge) {
private void syncDevices(Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncDevices [{}]", edge.getTenantId(), edge.getName());
try {
ListenableFuture<TimePageData<Device>> future =
deviceService.findDevicesByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE));
deviceService.findDevicesByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<Device>>() {
@Override
public void onSuccess(@Nullable TimePageData<Device> pageData) {
@ -194,6 +203,9 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
for (Device device : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ADDED, device.getId(), null);
}
if (pageData.hasNext()) {
syncDevices(edge, pageData.getNextPageLink());
}
}
}
@ -207,10 +219,10 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncAssets(Edge edge) {
private void syncAssets(Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncAssets [{}]", edge.getTenantId(), edge.getName());
try {
ListenableFuture<TimePageData<Asset>> future = assetService.findAssetsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE));
ListenableFuture<TimePageData<Asset>> future = assetService.findAssetsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<Asset>>() {
@Override
public void onSuccess(@Nullable TimePageData<Asset> pageData) {
@ -219,6 +231,9 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
for (Asset asset : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ASSET, EdgeEventActionType.ADDED, asset.getId(), null);
}
if (pageData.hasNext()) {
syncAssets(edge, pageData.getNextPageLink());
}
}
}
@ -232,10 +247,10 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncEntityViews(Edge edge) {
private void syncEntityViews(Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncEntityViews [{}]", edge.getTenantId(), edge.getName());
try {
ListenableFuture<TimePageData<EntityView>> future = entityViewService.findEntityViewsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE));
ListenableFuture<TimePageData<EntityView>> future = entityViewService.findEntityViewsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<EntityView>>() {
@Override
public void onSuccess(@Nullable TimePageData<EntityView> pageData) {
@ -244,6 +259,9 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
for (EntityView entityView : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ENTITY_VIEW, EdgeEventActionType.ADDED, entityView.getId(), null);
}
if (pageData.hasNext()) {
syncEntityViews(edge, pageData.getNextPageLink());
}
}
}
@ -257,10 +275,10 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncDashboards(Edge edge) {
private void syncDashboards(Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncDashboards [{}]", edge.getTenantId(), edge.getName());
try {
ListenableFuture<TimePageData<DashboardInfo>> future = dashboardService.findDashboardsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), new TimePageLink(Integer.MAX_VALUE));
ListenableFuture<TimePageData<DashboardInfo>> future = dashboardService.findDashboardsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<DashboardInfo>>() {
@Override
public void onSuccess(@Nullable TimePageData<DashboardInfo> pageData) {
@ -269,6 +287,9 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
for (DashboardInfo dashboardInfo : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DASHBOARD, EdgeEventActionType.ADDED, dashboardInfo.getId(), null);
}
if (pageData.hasNext()) {
syncDashboards(edge, pageData.getNextPageLink());
}
}
}
@ -285,18 +306,36 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private void syncUsers(Edge edge) {
log.trace("[{}] syncUsers [{}]", edge.getTenantId(), edge.getName());
try {
TextPageData<User> pageData = userService.findTenantAdmins(edge.getTenantId(), new TextPageLink(Integer.MAX_VALUE));
pushUsersToEdge(pageData, edge);
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null);
pageData = userService.findCustomerUsers(edge.getTenantId(), edge.getCustomerId(), new TextPageLink(Integer.MAX_VALUE));
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<User> pageData;
do {
pageData = userService.findTenantAdmins(edge.getTenantId(), pageLink);
pushUsersToEdge(pageData, edge);
}
syncCustomerUsers(edge);
if (pageData != null && pageData.hasNext()) {
pageLink = pageData.getNextPageLink();
}
} while (pageData != null && pageData.hasNext());
} catch (Exception 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);
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<User> pageData;
do {
pageData = userService.findCustomerUsers(edge.getTenantId(), edge.getCustomerId(), pageLink);
pushUsersToEdge(pageData, edge);
if (pageData != null && pageData.hasNext()) {
pageLink = pageData.getNextPageLink();
}
} while (pageData != null && pageData.hasNext());
}
}
private void syncWidgetsBundleAndWidgetTypes(Edge edge) {
log.trace("[{}] syncWidgetsBundleAndWidgetTypes [{}]", edge.getTenantId(), edge.getName());
List<WidgetsBundle> widgetsBundlesToPush = new ArrayList<>();
@ -426,7 +465,8 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
final EdgeEventType type = getEdgeQueueTypeByEntityType(entityId.getEntityType());
if (type != null) {
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>>() {
@Override
public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) {
@ -446,7 +486,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
entityData.put("kv", attributes);
entityData.put("scope", DataConstants.SERVER_SCOPE);
entityData.put("scope", scope);
JsonNode body = mapper.valueToTree(entityData);
log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body);
saveEdgeEvent(edge.getTenantId(),
@ -459,6 +499,11 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
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);
}
} else {
log.trace("[{}][{}] No attributes found for entity {} [{}]", edge.getTenantId(),
edge.getName(),
entityId.getEntityType(),
entityId.getId());
}
futureToSet.set(null);
}
@ -470,10 +515,8 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}, dbCallbackExecutorService);
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 {
log.warn("[{}] Type doesn't supported {}", edge.getTenantId(), entityId.getEntityType());
return Futures.immediateFuture(null);
}
}
@ -585,11 +628,11 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
private ListenableFuture<EdgeEvent> saveEdgeEvent(TenantId tenantId,
EdgeId edgeId,
EdgeEventType type,
EdgeEventActionType action,
EntityId entityId,
JsonNode body) {
EdgeId edgeId,
EdgeEventType type,
EdgeEventActionType action,
EntityId entityId,
JsonNode body) {
log.trace("Pushing edge event to edge queue. tenantId [{}], edgeId [{}], type [{}], action[{}], entityId [{}], body [{}]",
tenantId, edgeId, type, action, entityId, body);
@ -602,6 +645,18 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
edgeEvent.setEntityId(entityId.getId());
}
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.ObjectMapper;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
@ -107,6 +110,18 @@ public abstract class BaseProcessor {
edgeEvent.setEntityId(entityId.getId());
}
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;
}
}

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

@ -22,10 +22,12 @@ import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
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.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
@ -163,11 +165,31 @@ public class DefaultTbClusterService implements TbClusterService {
broadcast(new ComponentLifecycleMsg(tenantId, entityId, state));
}
@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) {
byte[] msgBytes = encodingService.encode(msg);
TbQueueProducer<TbProtoQueueMsg<ToRuleEngineNotificationMsg>> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer();
Set<String> tbRuleEngineServices = new HashSet<>(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE));
if (msg.getEntityId().getEntityType().equals(EntityType.TENANT)) {
boolean toCore = msg.getEntityId().getEntityType().equals(EntityType.TENANT) ||
msg.getEntityId().getEntityType().equals(EntityType.EDGE);
boolean toRuleEngine = !msg.getEntityId().getEntityType().equals(EntityType.EDGE);
if (toCore) {
TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer();
Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE);
for (String serviceId : tbCoreServices) {
@ -179,11 +201,13 @@ public class DefaultTbClusterService implements TbClusterService {
// No need to push notifications twice
tbRuleEngineServices.removeAll(tbCoreServices);
}
for (String serviceId : tbRuleEngineServices) {
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceId);
ToRuleEngineNotificationMsg toRuleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setComponentLifecycleMsg(ByteString.copyFrom(msgBytes)).build();
toRuleEngineProducer.send(tpi, new TbProtoQueueMsg<>(msg.getEntityId().getId(), toRuleEngineMsg), null);
toRuleEngineNfs.incrementAndGet();
if (toRuleEngine) {
for (String serviceId : tbRuleEngineServices) {
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceId);
ToRuleEngineNotificationMsg toRuleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setComponentLifecycleMsg(ByteString.copyFrom(msgBytes)).build();
toRuleEngineProducer.send(tpi, new TbProtoQueueMsg<>(msg.getEntityId().getId(), toRuleEngineMsg), null);
toRuleEngineNfs.incrementAndGet();
}
}
}

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

@ -203,18 +203,24 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
log.trace("[{}] Forwarding message to RPC service {}", id, toCoreNotification.getFromDeviceRpcResponse());
forwardToCoreRpcService(toCoreNotification.getFromDeviceRpcResponse(), callback);
} else if (toCoreNotification.getComponentLifecycleMsg() != null && !toCoreNotification.getComponentLifecycleMsg().isEmpty()) {
Optional<TbActorMsg> actorMsg = encodingService.decode(toCoreNotification.getComponentLifecycleMsg().toByteArray());
if (actorMsg.isPresent()) {
log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get());
actorContext.tellWithHighPriority(actorMsg.get());
}
callback.onSuccess();
forwardToAppActor(id, toCoreNotification.getComponentLifecycleMsg().toByteArray(), callback);
} else if (toCoreNotification.getEdgeEventUpdateMsg() != null && !toCoreNotification.getEdgeEventUpdateMsg().isEmpty()) {
forwardToAppActor(id, toCoreNotification.getEdgeEventUpdateMsg().toByteArray(), callback);
}
if (statsEnabled) {
stats.log(toCoreNotification);
}
}
private void forwardToAppActor(UUID id, byte[] msgBytes, TbCallback callback) {
Optional<TbActorMsg> actorMsg = encodingService.decode(msgBytes);
if (actorMsg.isPresent()) {
log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get());
actorContext.tellWithHighPriority(actorMsg.get());
}
callback.onSuccess();
}
private void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbCallback callback) {
RpcError error = proto.getError() > 0 ? RpcError.values()[proto.getError()] : null;
FromDeviceRpcResponse response = new FromDeviceRpcResponse(new UUID(proto.getRequestIdMSB(), proto.getRequestIdLSB())

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

@ -16,6 +16,7 @@
package org.thingsboard.server.service.queue;
import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
@ -49,4 +50,6 @@ public interface TbClusterService {
void onEntityStateChange(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state);
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 TbServiceInfoProvider serviceInfoProvider;
private final ConcurrentMap<UUID, Consumer<FromDeviceRpcResponse>> toDeviceRpcRequests = new ConcurrentHashMap<>();
private Optional<TbCoreDeviceRpcService> tbCoreRpcService;

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.WidgetsBundleUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.queue.TbClusterService;
import java.util.ArrayList;
import java.util.List;
@ -122,6 +123,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
@Autowired
private EdgeEventService edgeEventService;
@Autowired
private TbClusterService clusterService;
@Before
public void beforeTest() throws Exception {
loginSysAdmin();
@ -227,6 +231,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -847,6 +852,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -885,6 +891,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_DELETED, device.getId().getId(), EdgeEventType.DEVICE, deleteAttributesEntityData);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -910,6 +917,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.POST_ATTRIBUTES, device.getId().getId(), EdgeEventType.DEVICE, postAttributesEntityData);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -934,6 +942,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1);
clusterService.onEdgeEventUpdate(tenantId, edge.getId());
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
@ -1160,6 +1169,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
edgeImitator.sendUplinkMsg(uplinkMsgBuilder2.build());
edgeImitator.waitForResponses();
// Wait before device attributes saved to database before requesting them from controller
Thread.sleep(1000);
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));
@ -1302,18 +1312,25 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void sendAttributesRequest() throws Exception {
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);
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);
// Wait before device attributes saved to database before requesting them from edge
Thread.sleep(1000);
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
AttributesRequestMsg.Builder attributesRequestMsgBuilder = AttributesRequestMsg.newBuilder();
attributesRequestMsgBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits());
attributesRequestMsgBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits());
attributesRequestMsgBuilder.setEntityType(EntityType.DEVICE.name());
attributesRequestMsgBuilder.setScope(scope);
testAutoGeneratedCodeByProtobuf(attributesRequestMsgBuilder);
uplinkMsgBuilder.addAttributesRequestMsg(attributesRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
@ -1330,14 +1347,14 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB());
Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope());
Assert.assertEquals(scope, latestEntityDataMsg.getPostAttributeScope());
Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg());
TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg();
Assert.assertEquals(1, attributesUpdatedMsg.getKvCount());
TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0);
Assert.assertEquals("key1", keyValueProto.getKey());
Assert.assertEquals("value1", keyValueProto.getStringV());
Assert.assertEquals(expectedKey, keyValueProto.getKey());
Assert.assertEquals(expectedValue, keyValueProto.getStringV());
}
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.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@Slf4j
public class EdgeImitator {
@ -65,6 +67,8 @@ public class EdgeImitator {
private EdgeRpcClient edgeRpcClient;
private final Lock lock = new ReentrantLock();
private CountDownLatch messagesLatch;
private CountDownLatch responsesLatch;
private List<Class<? extends AbstractMessage>> ignoredTypes;
@ -74,7 +78,7 @@ public class EdgeImitator {
@Getter
private UserId userId;
@Getter
private List<AbstractMessage> downlinkMsgs;
private final List<AbstractMessage> downlinkMsgs;
public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException {
edgeRpcClient = new EdgeGrpcClient();
@ -241,7 +245,12 @@ public class EdgeImitator {
private ListenableFuture<Void> saveDownlinkMsg(AbstractMessage message) {
if (!ignoredTypes.contains(message.getClass())) {
downlinkMsgs.add(message);
try {
lock.lock();
downlinkMsgs.add(message);
} finally {
lock.unlock();
}
messagesLatch.countDown();
}
return Futures.immediateFuture(null);
@ -262,7 +271,14 @@ public class EdgeImitator {
}
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() {

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

@ -315,6 +315,7 @@ message AttributesRequestMsg {
int64 entityIdMSB = 1;
int64 entityIdLSB = 2;
string entityType = 3;
string scope = 4;
}
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.
*/
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

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

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

@ -98,6 +98,8 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId ";
public static final String INCORRECT_EDGE_ID = "Incorrect edgeId ";
private static final int DEFAULT_LIMIT = 100;
private RestTemplate restTemplate;
private static final String EDGE_LICENSE_SERVER_ENDPOINT = "https://license.thingsboard.io";
@ -460,8 +462,21 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
public ListenableFuture<List<EdgeId>> findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId) {
log.trace("[{}] Executing findRelatedEdgeIdsByEntityId [{}]", tenantId, entityId);
if (EntityType.TENANT.equals(entityId.getEntityType())) {
TextPageData<Edge> edgesByTenantId = findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE));
return Futures.immediateFuture(edgesByTenantId.getData().stream().map(IdBased::getId).collect(Collectors.toList()));
List<EdgeId> result = new ArrayList<>();
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<Edge> pageData;
do {
pageData = findEdgesByTenantId(tenantId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (Edge edge : pageData.getData()) {
result.add(edge.getId());
}
if (pageData.hasNext()) {
pageLink = pageData.getNextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
return Futures.immediateFuture(result);
} else {
switch (entityId.getEntityType()) {
case DEVICE:
@ -486,13 +501,23 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
if (userById == null) {
return Futures.immediateFuture(Collections.emptyList());
}
TextPageData<Edge> edges;
if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) {
edges = findEdgesByTenantId(tenantId, new TextPageLink(Integer.MAX_VALUE));
} else {
edges = findEdgesByTenantIdAndCustomerId(tenantId, new CustomerId(entityId.getId()), new TextPageLink(Integer.MAX_VALUE));
}
return convertToEdgeIds(Futures.immediateFuture(edges.getData()));
List<Edge> result = new ArrayList<>();
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<Edge> pageData;
do {
if (userById.getCustomerId() == null || userById.getCustomerId().isNullUid()) {
pageData = findEdgesByTenantId(tenantId, pageLink);
} 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 = pageData.getNextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
return convertToEdgeIds(Futures.immediateFuture(result));
default:
return Futures.immediateFuture(Collections.emptyList());
}

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

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

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

@ -115,11 +115,12 @@ public class TbMsgPushToEdgeNode implements TbNode {
@Override
public void onSuccess(@Nullable EdgeEvent event) {
ctx.tellNext(msg, SUCCESS);
ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId);
}
@Override
public void onFailure(Throwable th) {
log.error("Could not save edge event", th);
log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th);
ctx.tellFailure(msg, th);
}
}, ctx.getDbCallbackExecutor());

Loading…
Cancel
Save