|
|
|
@ -250,10 +250,17 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh |
|
|
|
checkActive(msg); |
|
|
|
EntityId entityId = msg.getOriginator(); |
|
|
|
TopicPartitionInfo tpi = systemContext.resolve(ServiceType.TB_RULE_ENGINE, msg.getQueueName(), tenantId, entityId); |
|
|
|
List<RuleNodeRelation> relations = nodeRoutes.get(originatorNodeId).stream() |
|
|
|
|
|
|
|
List<RuleNodeRelation> ruleNodeRelations = nodeRoutes.get(originatorNodeId); |
|
|
|
if (ruleNodeRelations == null) { // When unchecked, this will cause NullPointerException when rule node doesn't exist anymore
|
|
|
|
log.warn("[{}][{}][{}] No outbound relations (null). Probably rule node does not exist. Probably old message.", tenantId, entityId, msg.getId()); |
|
|
|
ruleNodeRelations = Collections.emptyList(); |
|
|
|
} |
|
|
|
|
|
|
|
List<RuleNodeRelation> relationsByTypes = ruleNodeRelations.stream() |
|
|
|
.filter(r -> contains(relationTypes, r.getType())) |
|
|
|
.collect(Collectors.toList()); |
|
|
|
int relationsCount = relations.size(); |
|
|
|
int relationsCount = relationsByTypes.size(); |
|
|
|
if (relationsCount == 0) { |
|
|
|
log.trace("[{}][{}][{}] No outbound relations to process", tenantId, entityId, msg.getId()); |
|
|
|
if (relationTypes.contains(TbRelationTypes.FAILURE)) { |
|
|
|
@ -268,14 +275,14 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh |
|
|
|
msg.getCallback().onSuccess(); |
|
|
|
} |
|
|
|
} else if (relationsCount == 1) { |
|
|
|
for (RuleNodeRelation relation : relations) { |
|
|
|
for (RuleNodeRelation relation : relationsByTypes) { |
|
|
|
log.trace("[{}][{}][{}] Pushing message to single target: [{}]", tenantId, entityId, msg.getId(), relation.getOut()); |
|
|
|
pushToTarget(tpi, msg, relation.getOut(), relation.getType()); |
|
|
|
} |
|
|
|
} else { |
|
|
|
MultipleTbQueueTbMsgCallbackWrapper callbackWrapper = new MultipleTbQueueTbMsgCallbackWrapper(relationsCount, msg.getCallback()); |
|
|
|
log.trace("[{}][{}][{}] Pushing message to multiple targets: [{}]", tenantId, entityId, msg.getId(), relations); |
|
|
|
for (RuleNodeRelation relation : relations) { |
|
|
|
log.trace("[{}][{}][{}] Pushing message to multiple targets: [{}]", tenantId, entityId, msg.getId(), relationsByTypes); |
|
|
|
for (RuleNodeRelation relation : relationsByTypes) { |
|
|
|
EntityId target = relation.getOut(); |
|
|
|
putToQueue(tpi, msg, callbackWrapper, target); |
|
|
|
} |
|
|
|
|