Browse Source

Singleton rule node implementation

pull/8412/head
Andrii Shvaika 4 years ago
parent
commit
b20057004b
  1. 1
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
  2. 7
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
  3. 5
      common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleNode.java

1
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java

@ -388,6 +388,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
private void pushMsgToNode(RuleNodeCtx nodeCtx, TbMsg msg, String fromRelationType) { private void pushMsgToNode(RuleNodeCtx nodeCtx, TbMsg msg, String fromRelationType) {
if (nodeCtx != null) { if (nodeCtx != null) {
//TODO: analyze that singleton flag and execute putToQueue if needed.
nodeCtx.getSelfActor().tell(new RuleChainToRuleNodeMsg(new DefaultTbContext(systemContext, ruleChainName, nodeCtx), msg, fromRelationType)); nodeCtx.getSelfActor().tell(new RuleChainToRuleNodeMsg(new DefaultTbContext(systemContext, ruleChainName, nodeCtx), msg, fromRelationType));
} else { } else {
log.error("[{}][{}] RuleNodeCtx is empty", entityId, ruleChainName); log.error("[{}][{}] RuleNodeCtx is empty", entityId, ruleChainName);

7
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java

@ -39,11 +39,10 @@ import org.thingsboard.server.common.stats.TbApiUsageReportClient;
public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNodeId> { public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNodeId> {
private final String ruleChainName; private final String ruleChainName;
private final TbActorRef self;
private final TbApiUsageReportClient apiUsageClient; private final TbApiUsageReportClient apiUsageClient;
private final DefaultTbContext defaultCtx;
private RuleNode ruleNode; private RuleNode ruleNode;
private TbNode tbNode; private TbNode tbNode;
private DefaultTbContext defaultCtx;
private RuleNodeInfo info; private RuleNodeInfo info;
RuleNodeActorMessageProcessor(TenantId tenantId, String ruleChainName, RuleNodeId ruleNodeId, ActorSystemContext systemContext RuleNodeActorMessageProcessor(TenantId tenantId, String ruleChainName, RuleNodeId ruleNodeId, ActorSystemContext systemContext
@ -51,7 +50,6 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
super(systemContext, tenantId, ruleNodeId); super(systemContext, tenantId, ruleNodeId);
this.apiUsageClient = systemContext.getApiUsageClient(); this.apiUsageClient = systemContext.getApiUsageClient();
this.ruleChainName = ruleChainName; this.ruleChainName = ruleChainName;
this.self = self;
this.ruleNode = systemContext.getRuleChainService().findRuleNodeById(tenantId, entityId); this.ruleNode = systemContext.getRuleChainService().findRuleNodeById(tenantId, entityId);
this.defaultCtx = new DefaultTbContext(systemContext, ruleChainName, new RuleNodeCtx(tenantId, parent, self, ruleNode)); this.defaultCtx = new DefaultTbContext(systemContext, ruleChainName, new RuleNodeCtx(tenantId, parent, self, ruleNode));
this.info = new RuleNodeInfo(ruleNodeId, ruleChainName, ruleNode != null ? ruleNode.getName() : "Unknown"); this.info = new RuleNodeInfo(ruleNodeId, ruleChainName, ruleNode != null ? ruleNode.getName() : "Unknown");
@ -59,6 +57,7 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
@Override @Override
public void start(TbActorCtx context) throws Exception { public void start(TbActorCtx context) throws Exception {
//TODO: do not start the node if singleton
tbNode = initComponent(ruleNode); tbNode = initComponent(ruleNode);
if (tbNode != null) { if (tbNode != null) {
state = ComponentLifecycleState.ACTIVE; state = ComponentLifecycleState.ACTIVE;
@ -95,12 +94,14 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
@Override @Override
public void onPartitionChangeMsg(PartitionChangeMsg msg) { public void onPartitionChangeMsg(PartitionChangeMsg msg) {
//TODO: start the node
if (tbNode != null) { if (tbNode != null) {
tbNode.onPartitionChangeMsg(defaultCtx, msg); tbNode.onPartitionChangeMsg(defaultCtx, msg);
} }
} }
public void onRuleToSelfMsg(RuleNodeToSelfMsg msg) throws Exception { public void onRuleToSelfMsg(RuleNodeToSelfMsg msg) throws Exception {
// TODO: check that the rule node is singleton and use putToQueue
checkComponentStateActive(msg.getMsg()); checkComponentStateActive(msg.getMsg());
TbMsg tbMsg = msg.getMsg(); TbMsg tbMsg = msg.getMsg();
int ruleNodeCount = tbMsg.getAndIncrementRuleNodeCounter(); int ruleNodeCount = tbMsg.getAndIncrementRuleNodeCounter();

5
common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleNode.java

@ -48,7 +48,9 @@ public class RuleNode extends SearchTextBasedWithAdditionalInfo<RuleNodeId> impl
private String name; private String name;
@ApiModelProperty(position = 6, value = "Enable/disable debug. ", example = "false") @ApiModelProperty(position = 6, value = "Enable/disable debug. ", example = "false")
private boolean debugMode; private boolean debugMode;
@ApiModelProperty(position = 7, value = "JSON with the rule node configuration. Structure depends on the rule node implementation.", dataType = "com.fasterxml.jackson.databind.JsonNode") @ApiModelProperty(position = 7, value = "Enable/disable singleton mode. ", example = "false")
private boolean singletonMode;
@ApiModelProperty(position = 8, value = "JSON with the rule node configuration. Structure depends on the rule node implementation.", dataType = "com.fasterxml.jackson.databind.JsonNode")
private transient JsonNode configuration; private transient JsonNode configuration;
@JsonIgnore @JsonIgnore
private byte[] configurationBytes; private byte[] configurationBytes;
@ -69,6 +71,7 @@ public class RuleNode extends SearchTextBasedWithAdditionalInfo<RuleNodeId> impl
this.type = ruleNode.getType(); this.type = ruleNode.getType();
this.name = ruleNode.getName(); this.name = ruleNode.getName();
this.debugMode = ruleNode.isDebugMode(); this.debugMode = ruleNode.isDebugMode();
this.singletonMode = ruleNode.isSingletonMode();
this.setConfiguration(ruleNode.getConfiguration()); this.setConfiguration(ruleNode.getConfiguration());
this.externalId = ruleNode.getExternalId(); this.externalId = ruleNode.getExternalId();
} }

Loading…
Cancel
Save