diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index 38d6f4bccb..609c087a51 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -167,6 +167,11 @@ class DefaultTbContext implements TbContext { } private void enqueue(TopicPartitionInfo tpi, TbMsg tbMsg, Consumer onFailure, Runnable onSuccess) { + if (!tbMsg.isValid()) { + log.trace("[{}] Skip invalid message: {}", getTenantId(), tbMsg); + onFailure.accept(new IllegalArgumentException("Source message is no longer valid!")); + return; + } TransportProtos.ToRuleEngineMsg msg = TransportProtos.ToRuleEngineMsg.newBuilder() .setTenantIdMSB(getTenantId().getId().getMostSignificantBits()) .setTenantIdLSB(getTenantId().getId().getLeastSignificantBits()) @@ -235,6 +240,11 @@ class DefaultTbContext implements TbContext { } private void enqueueForTellNext(TopicPartitionInfo tpi, String queueName, TbMsg source, Set relationTypes, String failureMessage, Runnable onSuccess, Consumer onFailure) { + if (!source.isValid()) { + log.trace("[{}] Skip invalid message: {}", getTenantId(), source); + onFailure.accept(new IllegalArgumentException("Source message is no longer valid!")); + return; + } RuleChainId ruleChainId = nodeCtx.getSelf().getRuleChainId(); RuleNodeId ruleNodeId = nodeCtx.getSelf().getId(); TbMsg tbMsg = TbMsg.newMsg(source, queueName, ruleChainId, ruleNodeId); diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java index 88a4f6a202..6f6ea8fcec 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java @@ -200,17 +200,20 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor relationTypes, String failureMessage) { try { - checkActive(msg); + checkComponentStateActive(msg); EntityId entityId = msg.getOriginator(); TopicPartitionInfo tpi = systemContext.resolve(ServiceType.TB_RULE_ENGINE, msg.getQueueName(), tenantId, entityId); diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActor.java index 8fef00c49c..60f462f4d4 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActor.java @@ -27,6 +27,7 @@ import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbActorMsg; +import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; @@ -86,12 +87,19 @@ public class RuleNodeActor extends ComponentActor extends Abstract schedulePeriodicMsgWithDelay(context, new StatsPersistTick(), statsPersistFrequency, statsPersistFrequency); } - protected void checkActive(TbMsg tbMsg) throws RuleNodeException { + protected boolean checkMsgValid(TbMsg tbMsg) { + var valid = tbMsg.isValid(); + if (!valid) { + if (log.isTraceEnabled()) { + log.trace("Skip processing of message: {} because it is no longer valid!", tbMsg); + } + } + return valid; + } + + protected void checkComponentStateActive(TbMsg tbMsg) throws RuleNodeException { if (state != ComponentLifecycleState.ACTIVE) { log.debug("Component is not active. Current state [{}] for processor [{}][{}] tenant [{}]", state, entityId.getEntityType(), entityId, tenantId); RuleNodeException ruleNodeException = getInactiveException(); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java index eff1ecaf86..e3ce0927ef 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackCallback.java @@ -66,6 +66,11 @@ public class TbMsgPackCallback implements TbMsgCallback { ctx.onFailure(tenantId, id, e); } + @Override + public boolean isMsgValid() { + return !ctx.isComplete(); + } + @Override public void onProcessingStart(RuleNodeInfo ruleNodeInfo) { log.trace("[{}] ON PROCESSING START: {}", id, ruleNodeInfo); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java index d7a064c4ca..62285a411c 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java @@ -53,6 +53,8 @@ public class TbMsgPackProcessingContext { private final ConcurrentMap exceptionsMap = new ConcurrentHashMap<>(); private final ConcurrentMap lastRuleNodeMap = new ConcurrentHashMap<>(); + @Getter + private volatile boolean complete = false; public TbMsgPackProcessingContext(String queueName, TbRuleEngineSubmitStrategy submitStrategy) { this.queueName = queueName; @@ -149,6 +151,7 @@ public class TbMsgPackProcessingContext { } public void cleanup() { + complete = true; pendingMap.clear(); successMap.clear(); failedMap.clear(); diff --git a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java index 90a0404751..eb79bc9d35 100644 --- a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java @@ -144,6 +144,7 @@ public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRule Thread.sleep(1000); TbMsgCallback tbMsgCallback = Mockito.mock(TbMsgCallback.class); + Mockito.when(tbMsgCallback.isMsgValid()).thenReturn(true); TbMsg tbMsg = TbMsg.newMsg("CUSTOM", device.getId(), new TbMsgMetaData(), "{}", tbMsgCallback); QueueToRuleEngineMsg qMsg = new QueueToRuleEngineMsg(savedTenant.getId(), tbMsg, null, null); // Pushing Message to the system @@ -256,6 +257,7 @@ public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRule Thread.sleep(1000); TbMsgCallback tbMsgCallback = Mockito.mock(TbMsgCallback.class); + Mockito.when(tbMsgCallback.isMsgValid()).thenReturn(true); TbMsg tbMsg = TbMsg.newMsg("CUSTOM", device.getId(), new TbMsgMetaData(), "{}", tbMsgCallback); QueueToRuleEngineMsg qMsg = new QueueToRuleEngineMsg(savedTenant.getId(), tbMsg, null, null); // Pushing Message to the system diff --git a/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java b/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java index ccdf788634..4a92321e91 100644 --- a/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java @@ -140,6 +140,7 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac Thread.sleep(1000); TbMsgCallback tbMsgCallback = Mockito.mock(TbMsgCallback.class); + Mockito.when(tbMsgCallback.isMsgValid()).thenReturn(true); TbMsg tbMsg = TbMsg.newMsg("CUSTOM", device.getId(), new TbMsgMetaData(), "{}", tbMsgCallback); QueueToRuleEngineMsg qMsg = new QueueToRuleEngineMsg(savedTenant.getId(), tbMsg, null, null); // Pushing Message to the system diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java index 63c148b951..7c77d8a06b 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java @@ -269,7 +269,7 @@ public final class TbMsg implements Serializable { } public TbMsgCallback getCallback() { - //May be null in case of deserialization; + // May be null in case of deserialization; if (callback != null) { return callback; } else { @@ -288,4 +288,12 @@ public final class TbMsg implements Serializable { public TbMsgProcessingStackItem popFormStack() { return ctx.pop(); } + + /** + * Checks if the message is still valid for processing. May be invalid if the message pack is timed-out or canceled. + * @return 'true' if message is valid for processing, 'false' otherwise. + */ + public boolean isValid() { + return getCallback().isMsgValid(); + } } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TbMsgCallback.java b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TbMsgCallback.java index fd62df83c9..2bab1e8340 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TbMsgCallback.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TbMsgCallback.java @@ -17,6 +17,9 @@ package org.thingsboard.server.common.msg.queue; import org.thingsboard.server.common.data.id.RuleNodeId; +/** + * Should be renamed to TbMsgPackContext, but this can't be changed due to backward-compatibility. + */ public interface TbMsgCallback { TbMsgCallback EMPTY = new TbMsgCallback() { @@ -36,11 +39,20 @@ public interface TbMsgCallback { void onFailure(RuleEngineException e); + /** + * Returns 'true' if rule engine is expecting the message to be processed, 'false' otherwise. + * message may no longer be valid, if the message pack is already expired/canceled/failed. + * + * @return 'true' if rule engine is expecting the message to be processed, 'false' otherwise. + */ + default boolean isMsgValid() { + return true; + } + default void onProcessingStart(RuleNodeInfo ruleNodeInfo) { } default void onProcessingEnd(RuleNodeId ruleNodeId) { } - }