Browse Source

Handle asset removal from edge event in AssetEdgeProcessor

pull/14447/head
Nikita Mazurenko 10 months ago
parent
commit
14171e30a1
  1. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java
  2. 5
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessor.java
  3. 18
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/BaseAssetProcessor.java
  4. 16
      application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java

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

@ -363,7 +363,7 @@ public abstract class BaseEdgeProcessor implements EdgeProcessor {
pushEntityEventToRuleEngine(tenantId, entity.getId(), customerId, msgType, entityAsString, tbMsgMetaData);
} catch (Exception e) {
log.warn("[{}][{}] Failed to push entity of type {} action to rule engine: {}", tenantId, entity.getId(), entity.getId().getEntityType(), msgType.name(), e);
log.warn("[{}][{}] Failed to push entity action for {} to rule engine: {}", tenantId, entity.getId(), entity.getId().getEntityType(), msgType.name(), e);
}
}

5
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessor.java

@ -61,10 +61,7 @@ public class AssetEdgeProcessor extends BaseAssetProcessor implements AssetProce
saveOrUpdateAsset(tenantId, assetId, assetUpdateMsg, edge);
return Futures.immediateFuture(null);
case ENTITY_DELETED_RPC_MESSAGE:
Asset assetToDelete = edgeCtx.getAssetService().findAssetById(tenantId, assetId);
if (assetToDelete != null) {
edgeCtx.getAssetService().unassignAssetFromEdge(tenantId, assetId, edge.getId());
}
deleteAsset(tenantId, edge, assetId);
return Futures.immediateFuture(null);
case UNRECOGNIZED:
default:

18
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/BaseAssetProcessor.java

@ -21,9 +21,11 @@ import org.springframework.data.util.Pair;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
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.gen.edge.v1.AssetUpdateMsg;
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor;
@ -77,4 +79,20 @@ public abstract class BaseAssetProcessor extends BaseEdgeProcessor {
protected abstract void setCustomerId(TenantId tenantId, CustomerId customerId, Asset asset, AssetUpdateMsg assetUpdateMsg);
protected void deleteAsset(TenantId tenantId, AssetId assetId) {
Asset assetById = edgeCtx.getAssetService().findAssetById(tenantId, assetId);
if (assetById != null) {
edgeCtx.getAssetService().deleteAsset(tenantId, assetId);
pushEntityEventToRuleEngine(tenantId, null, assetById, TbMsgType.ENTITY_DELETED);
}
}
protected void deleteAsset(TenantId tenantId, Edge edge, AssetId assetId) {
Asset assetById = edgeCtx.getAssetService().findAssetById(tenantId, assetId);
if (assetById != null) {
edgeCtx.getAssetService().deleteAsset(tenantId, assetId);
pushEntityEventToRuleEngine(tenantId, edge, assetById, TbMsgType.ENTITY_DELETED);
}
}
}

16
application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java

@ -15,7 +15,6 @@
*/
package org.thingsboard.server.edge;
import com.fasterxml.jackson.core.type.TypeReference;
import com.google.protobuf.AbstractMessage;
import org.junit.Assert;
import org.junit.Test;
@ -28,8 +27,6 @@ import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
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.gen.edge.v1.AssetProfileUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg;
@ -37,10 +34,11 @@ import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.gen.edge.v1.UplinkMsg;
import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg;
import java.util.List;
import java.util.Optional;
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;
@DaoSqlTest
@ -277,12 +275,10 @@ public class AssetEdgeTest extends AbstractEdgeTest {
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(upLinkMsgBuilder.build());
Assert.assertTrue(edgeImitator.waitForResponses());
AssetInfo assetInfo = doGet("/api/asset/info/" + savedAsset.getUuidId(), AssetInfo.class);
Assert.assertNotNull(assetInfo);
List<AssetInfo> edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getUuidId() + "/assets?",
new TypeReference<PageData<AssetInfo>>() {
}, new PageLink(100)).getData();
Assert.assertFalse(edgeAssets.contains(assetInfo));
await().atMost(30, TimeUnit.SECONDS).untilAsserted(() ->
doGet("/api/asset/info/" + savedAsset.getUuidId(), AssetInfo.class, status().isNotFound())
);
}
private Asset saveAssetOnCloudAndVerifyDeliveryToEdge() throws Exception {

Loading…
Cancel
Save