Browse Source

Handle entity view removal from edge event in EntityViewEdgeProcessor

pull/14447/head
Nikita Mazurenko 10 months ago
parent
commit
394075c9df
  1. 13
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/entityview/BaseEntityViewProcessor.java
  2. 34
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/entityview/EntityViewEdgeProcessor.java
  3. 15
      application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java

13
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/entityview/BaseEntityViewProcessor.java

@ -21,9 +21,11 @@ import org.springframework.data.util.Pair;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg;
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor;
@ -69,4 +71,15 @@ public abstract class BaseEntityViewProcessor extends BaseEdgeProcessor {
protected abstract void setCustomerId(TenantId tenantId, CustomerId customerId, EntityView entityView, EntityViewUpdateMsg entityViewUpdateMsg); protected abstract void setCustomerId(TenantId tenantId, CustomerId customerId, EntityView entityView, EntityViewUpdateMsg entityViewUpdateMsg);
protected void deleteEntityView(TenantId tenantId, EntityViewId entityViewId) {
deleteEntityView(tenantId, null, entityViewId);
}
protected void deleteEntityView(TenantId tenantId, Edge edge, EntityViewId entityViewId) {
EntityView entityViewById = edgeCtx.getEntityViewService().findEntityViewById(tenantId, entityViewId);
if (entityViewById != null) {
edgeCtx.getEntityViewService().deleteEntityView(tenantId, entityViewId);
pushEntityEventToRuleEngine(tenantId, edge, entityViewById, TbMsgType.ENTITY_DELETED);
}
}
} }

34
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/entityview/EntityViewEdgeProcessor.java

@ -54,21 +54,17 @@ public class EntityViewEdgeProcessor extends BaseEntityViewProcessor implements
try { try {
edgeSynchronizationManager.getEdgeId().set(edge.getId()); edgeSynchronizationManager.getEdgeId().set(edge.getId());
switch (entityViewUpdateMsg.getMsgType()) { return switch (entityViewUpdateMsg.getMsgType()) {
case ENTITY_CREATED_RPC_MESSAGE: case ENTITY_CREATED_RPC_MESSAGE, ENTITY_UPDATED_RPC_MESSAGE -> {
case ENTITY_UPDATED_RPC_MESSAGE:
saveOrUpdateEntityView(tenantId, entityViewId, entityViewUpdateMsg, edge); saveOrUpdateEntityView(tenantId, entityViewId, entityViewUpdateMsg, edge);
return Futures.immediateFuture(null); yield Futures.immediateFuture(null);
case ENTITY_DELETED_RPC_MESSAGE: }
EntityView entityViewToDelete = edgeCtx.getEntityViewService().findEntityViewById(tenantId, entityViewId); case ENTITY_DELETED_RPC_MESSAGE -> {
if (entityViewToDelete != null) { deleteEntityView(tenantId, entityViewId);
edgeCtx.getEntityViewService().unassignEntityViewFromEdge(tenantId, entityViewId, edge.getId()); yield Futures.immediateFuture(null);
} }
return Futures.immediateFuture(null); default -> handleUnsupportedMsgType(entityViewUpdateMsg.getMsgType());
case UNRECOGNIZED: };
default:
return handleUnsupportedMsgType(entityViewUpdateMsg.getMsgType());
}
} catch (DataValidationException e) { } catch (DataValidationException e) {
if (e.getMessage().contains("limit reached")) { if (e.getMessage().contains("limit reached")) {
log.warn("[{}] Number of allowed entity views violated {}", tenantId, entityViewUpdateMsg, e); log.warn("[{}] Number of allowed entity views violated {}", tenantId, entityViewUpdateMsg, e);
@ -96,14 +92,8 @@ public class EntityViewEdgeProcessor extends BaseEntityViewProcessor implements
} }
private void pushEntityViewCreatedEventToRuleEngine(TenantId tenantId, Edge edge, EntityViewId entityViewId) { private void pushEntityViewCreatedEventToRuleEngine(TenantId tenantId, Edge edge, EntityViewId entityViewId) {
try { EntityView entityView = edgeCtx.getEntityViewService().findEntityViewById(tenantId, entityViewId);
EntityView entityView = edgeCtx.getEntityViewService().findEntityViewById(tenantId, entityViewId); pushEntityEventToRuleEngine(tenantId, edge, entityView, TbMsgType.ENTITY_CREATED);
String entityViewAsString = JacksonUtil.toString(entityView);
TbMsgMetaData msgMetaData = getEdgeActionTbMsgMetaData(edge, entityView.getCustomerId());
pushEntityEventToRuleEngine(tenantId, entityViewId, entityView.getCustomerId(), TbMsgType.ENTITY_CREATED, entityViewAsString, msgMetaData);
} catch (Exception e) {
log.warn("[{}][{}] Failed to push entity view action to rule engine: {}", tenantId, entityViewId, TbMsgType.ENTITY_CREATED.name(), e);
}
} }
@Override @Override

15
application/src/test/java/org/thingsboard/server/edge/EntityViewEdgeTest.java

@ -15,7 +15,6 @@
*/ */
package org.thingsboard.server.edge; package org.thingsboard.server.edge;
import com.fasterxml.jackson.core.type.TypeReference;
import com.google.protobuf.AbstractMessage; import com.google.protobuf.AbstractMessage;
import com.google.protobuf.InvalidProtocolBufferException; import com.google.protobuf.InvalidProtocolBufferException;
import org.junit.Assert; import org.junit.Assert;
@ -31,8 +30,6 @@ import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg;
import org.thingsboard.server.gen.edge.v1.EntityViewsRequestMsg; import org.thingsboard.server.gen.edge.v1.EntityViewsRequestMsg;
@ -40,10 +37,11 @@ import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UplinkMsg;
import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg;
import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.TimeUnit;
import static org.awaitility.Awaitility.await;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@DaoSqlTest @DaoSqlTest
@ -256,12 +254,9 @@ public class EntityViewEdgeTest extends AbstractEdgeTest {
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(upLinkMsgBuilder.build()); edgeImitator.sendUplinkMsg(upLinkMsgBuilder.build());
Assert.assertTrue(edgeImitator.waitForResponses()); Assert.assertTrue(edgeImitator.waitForResponses());
EntityViewInfo entityViewInfo = doGet("/api/entityView/info/" + savedEntityView.getUuidId(), EntityViewInfo.class); await().atMost(30, TimeUnit.SECONDS).untilAsserted(() ->
Assert.assertNotNull(entityViewInfo); doGet("/api/entityView/info/" + savedEntityView.getUuidId(), EntityViewInfo.class, status().isNotFound())
List<EntityViewInfo> edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getUuidId() + "/entityViews?", );
new TypeReference<PageData<EntityViewInfo>>() {
}, new PageLink(100)).getData();
Assert.assertFalse(edgeAssets.contains(entityViewInfo));
} }
private void verifyEntityViewUpdateMsg(EntityView entityView, Device device) throws InvalidProtocolBufferException { private void verifyEntityViewUpdateMsg(EntityView entityView, Device device) throws InvalidProtocolBufferException {

Loading…
Cancel
Save