Browse Source

Tech depts. Added support for Edge entity on push to edge node

pull/3693/head
Volodymyr Babak 6 years ago
parent
commit
0c7a51242d
  1. 17
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  2. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 20
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java
  4. 5
      common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java
  5. 8
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  6. 3
      dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java
  7. 16
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
  8. 104
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

17
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;
@ -445,7 +446,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
@Override @Override
public void onSuccess(@Nullable Alarm alarm) { public void onSuccess(@Nullable Alarm alarm) {
if (alarm != null) { if (alarm != null) {
EdgeEventType type = getEdgeQueueTypeByEntityType(alarm.getOriginator().getEntityType()); EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType());
if (type != null) { if (type != null) {
ListenableFuture<List<EdgeId>> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator()); ListenableFuture<List<EdgeId>> relatedEdgeIdsByEntityIdFuture = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator());
Futures.addCallback(relatedEdgeIdsByEntityIdFuture, new FutureCallback<List<EdgeId>>() { Futures.addCallback(relatedEdgeIdsByEntityIdFuture, new FutureCallback<List<EdgeId>>() {
@ -518,20 +519,6 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
} }
} }
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:
log.debug("Unsupported entity type: [{}]", entityType);
return null;
}
}
} }

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

@ -431,6 +431,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) {

20
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;
@ -146,7 +147,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
@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);
@ -157,7 +158,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
syncEntityViews(edge, new TimePageLink(DEFAULT_LIMIT)); syncEntityViews(edge, new TimePageLink(DEFAULT_LIMIT));
syncDashboards(edge, new TimePageLink(DEFAULT_LIMIT)); syncDashboards(edge, new TimePageLink(DEFAULT_LIMIT));
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during sync process", e); log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), e);
} }
} }
@ -461,7 +462,7 @@ 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();
String scope = attributesRequestMsg.getScope(); String scope = attributesRequestMsg.getScope();
@ -520,19 +521,6 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
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);

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

@ -5,7 +5,7 @@
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
@ -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;
} }
} }

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

@ -469,12 +469,16 @@ 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())) {
List<EdgeId> result = new ArrayList<>(); List<EdgeId> result = new ArrayList<>();
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT); TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<Edge> pageData; TextPageData<Edge> pageData;
do { do {
pageData = findEdgesByTenantId(tenantId, pageLink); 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()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
for (Edge edge : pageData.getData()) { for (Edge edge : pageData.getData()) {
result.add(edge.getId()); result.add(edge.getId());

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

@ -90,6 +90,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) {
@ -97,7 +98,7 @@ public abstract class AbstractEntityService {
} }
} }
} catch (ExecutionException | InterruptedException e) { } catch (ExecutionException | InterruptedException 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();
@ -131,13 +131,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!.");
} }
} }
}; };

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

@ -5,7 +5,7 @@
* you may not use this file except in compliance with the License. * you may not use this file except in compliance with the License.
* You may obtain a copy of the License at * You may obtain a copy of the License at
* *
* http://www.apache.org/licenses/LICENSE-2.0 * http://www.apache.org/licenses/LICENSE-2.0
* *
* Unless required by applicable law or agreed to in writing, software * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * distributed under the License is distributed on an "AS IS" BASIS,
@ -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,48 +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.onEdgeEventUpdate(ctx.getTenantId(), edgeId); }
} }, ctx.getDbCallbackExecutor());
}
@Override } catch (JsonProcessingException e) {
public void onFailure(Throwable th) { log.error("Failed to build edge event", e);
log.warn("[{}] Can't save edge event [{}] for edge [{}]", ctx.getTenantId().getId(), edgeEvent, edgeId.getId(), th); ctx.tellFailure(msg, e);
ctx.tellFailure(msg, th); }
} } else {
}, ctx.getDbCallbackExecutor()); 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 {
@ -221,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;

Loading…
Cancel
Save