Browse Source

Merge pull request #11588 from volodymyr-babak/edge-rule-chain-cache-fix

Fixed issue with Edge Root Rule Chain cache for related edges
pull/11617/head
Viacheslav Klimov 2 years ago
committed by GitHub
parent
commit
3fdf93d3bf
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 13
      application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/edge/RelatedEdgesSourcingListener.java
  3. 25
      application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
  4. 12
      application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java
  5. 4
      dao/src/main/java/org/thingsboard/server/dao/edge/BaseRelatedEdgesService.java
  6. 6
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

13
application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java

@ -34,8 +34,10 @@ import org.thingsboard.server.common.data.alarm.AlarmComment;
import org.thingsboard.server.common.data.alarm.EntityAlarm;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.domain.Domain;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
@ -134,6 +136,17 @@ public class EdgeEventSourcingListener {
return;
}
try {
if (event.getEntityId().getEntityType().equals(EntityType.RULE_CHAIN) && event.getEdgeId() != null && event.getActionType().equals(ActionType.ASSIGNED_TO_EDGE)) {
try {
Edge edge = JacksonUtil.fromString(event.getBody(), Edge.class);
if (edge != null && new RuleChainId(event.getEntityId().getId()).equals(edge.getRootRuleChainId())) {
log.trace("[{}] skipping ASSIGNED_TO_EDGE event of RULE_CHAIN entity in case Edge Root Rule Chain: {}", event.getTenantId(), event);
return;
}
} catch (Exception ignored) {
return;
}
}
log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event);
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), event.getEdgeId(), event.getEntityId(),
event.getBody(), null, EdgeUtils.getEdgeEventActionTypeByActionType(event.getActionType()),

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

@ -55,6 +55,7 @@ public class RelatedEdgesSourcingListener {
@TransactionalEventListener(fallbackExecution = true)
public void handleEvent(ActionEntityEvent<?> event) {
executorService.submit(() -> {
log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event);
try {
switch (event.getActionType()) {
case ASSIGNED_TO_EDGE, UNASSIGNED_FROM_EDGE ->
@ -69,6 +70,7 @@ public class RelatedEdgesSourcingListener {
@TransactionalEventListener(fallbackExecution = true)
public void handleEvent(DeleteEntityEvent<?> event) {
executorService.submit(() -> {
log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event);
try {
relatedEdgesService.publishRelatedEdgeIdsEvictEvent(event.getTenantId(), event.getEntityId());
} catch (Exception e) {

25
application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java

@ -103,6 +103,7 @@ import org.thingsboard.server.gen.edge.v1.UserUpdateMsg;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import java.util.Random;
import java.util.TreeMap;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
@ -124,6 +125,8 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
protected EdgeImitator edgeImitator;
protected Edge edge;
private Random random = new Random();
@Autowired
protected EdgeEventService edgeEventService;
@ -163,15 +166,10 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
}
private RuleChainId getEdgeRootRuleChainId() throws Exception {
List<RuleChain> edgeRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getUuidId() + "/ruleChains?",
new TypeReference<PageData<RuleChain>>() {
}, new PageLink(100)).getData();
for (RuleChain edgeRuleChain : edgeRuleChains) {
if (edgeRuleChain.isRoot()) {
return edgeRuleChain.getId();
}
}
throw new RuntimeException("Root rule chain not found");
return doGetTypedWithPageLink("/api/ruleChains?type={type}&", new TypeReference<PageData<RuleChain>>() {},
new PageLink(100, 0, "Edge Root Rule Chain"),
"EDGE")
.getData().get(0).getId();
}
@After
@ -196,6 +194,8 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
Asset savedAsset = saveAsset("Edge Asset 1");
updateRootRuleChainMetadata();
edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class);
doPost("/api/edge/" + edge.getUuidId()
@ -207,6 +207,13 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
TimeUnit.MILLISECONDS.sleep(1000);
}
protected void updateRootRuleChainMetadata() throws Exception {
RuleChainId rootRuleChainId = getEdgeRootRuleChainId();
RuleChainMetaData rootRuleChainMetadata = doGet("/api/ruleChain/" + rootRuleChainId.getId().toString() + "/metadata", RuleChainMetaData.class);
rootRuleChainMetadata.getNodes().forEach(n -> n.setDebugMode(random.nextBoolean()));
doPost("/api/ruleChain/metadata", rootRuleChainMetadata, RuleChainMetaData.class);
}
protected void extendDeviceProfileData(DeviceProfile deviceProfile) {
DeviceProfileData profileData = deviceProfile.getProfileData();
List<DeviceProfileAlarm> alarms = new ArrayList<>();

12
application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java

@ -185,6 +185,18 @@ public class RuleChainEdgeTest extends AbstractEdgeTest {
return doPost("/api/ruleChain/metadata", ruleChainMetaData, RuleChainMetaData.class);
}
@Test
public void testUpdateRootRuleChain() throws Exception {
edgeImitator.expectMessageAmount(2);
updateRootRuleChainMetadata();
Assert.assertTrue(edgeImitator.waitForMessages());
Optional<RuleChainUpdateMsg> ruleChainUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainUpdateMsg.class);
Assert.assertTrue(ruleChainUpdateMsgOpt.isPresent());
Optional<RuleChainMetadataUpdateMsg> ruleChainMetadataUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainMetadataUpdateMsg.class);
Assert.assertTrue(ruleChainMetadataUpdateMsgOpt.isPresent());
}
@Test
public void testSetRootRuleChain() throws Exception {
// create rule chain

4
dao/src/main/java/org/thingsboard/server/dao/edge/BaseRelatedEdgesService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.edge;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
@ -30,6 +31,7 @@ import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.entity.AbstractCachedEntityService;
@Service
@Slf4j
public class BaseRelatedEdgesService extends AbstractCachedEntityService<RelatedEdgesCacheKey, RelatedEdgesCacheValue, RelatedEdgesEvictEvent> implements RelatedEdgesService {
public static final int RELATED_EDGES_CACHE_ITEMS = 1000;
@ -47,6 +49,7 @@ public class BaseRelatedEdgesService extends AbstractCachedEntityService<Related
@Override
public PageData<EdgeId> findEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) {
log.trace("Executing findEdgeIdsByEntityId, tenantId [{}], entityId [{}], pageLink [{}]", tenantId, entityId, pageLink);
if (!pageLink.equals(FIRST_PAGE)) {
return edgeService.findEdgeIdsByTenantIdAndEntityId(tenantId, entityId, pageLink);
}
@ -56,6 +59,7 @@ public class BaseRelatedEdgesService extends AbstractCachedEntityService<Related
@Override
public void publishRelatedEdgeIdsEvictEvent(TenantId tenantId, EntityId entityId) {
log.trace("Executing publishRelatedEdgeIdsEvictEvent, tenantId [{}], entityId [{}]", tenantId, entityId);
publishEvictEvent(new RelatedEdgesEvictEvent(tenantId, entityId));
}

6
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -640,10 +640,8 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
log.warn("[{}] Failed to create ruleChain relation. Edge Id: [{}]", ruleChainId, edgeId);
throw new RuntimeException(e);
}
if (!ruleChainId.equals(edge.getRootRuleChainId())) {
eventPublisher.publishEvent(ActionEntityEvent.builder().tenantId(tenantId).edgeId(edgeId).entityId(ruleChainId)
.actionType(ActionType.ASSIGNED_TO_EDGE).build());
}
eventPublisher.publishEvent(ActionEntityEvent.builder().tenantId(tenantId).edgeId(edgeId).entityId(ruleChainId)
.actionType(ActionType.ASSIGNED_TO_EDGE).body(JacksonUtil.toString(edge)).build());
return ruleChain;
}

Loading…
Cancel
Save