Browse Source

TB-73: Implementation

pull/218/head
Andrew Shvayka 9 years ago
parent
commit
ae4a71346f
  1. 55
      application/src/main/java/org/thingsboard/server/actors/rule/RuleActorMessageProcessor.java
  2. 7
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleService.java

55
application/src/main/java/org/thingsboard/server/actors/rule/RuleActorMessageProcessor.java

@ -18,6 +18,7 @@ package org.thingsboard.server.actors.rule;
import java.util.*; import java.util.*;
import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.JsonProcessingException;
import org.springframework.util.StringUtils;
import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.plugin.RuleToPluginMsgWrapper; import org.thingsboard.server.actors.plugin.RuleToPluginMsgWrapper;
import org.thingsboard.server.actors.shared.ComponentMsgProcessor; import org.thingsboard.server.actors.shared.ComponentMsgProcessor;
@ -113,8 +114,9 @@ class RuleActorMessageProcessor extends ComponentMsgProcessor<RuleId> {
} }
private void initAction() throws Exception { private void initAction() throws Exception {
JsonNode actionMd = ruleMd.getAction(); if (ruleMd.getAction() != null && !ruleMd.getAction().isNull()) {
action = initComponent(actionMd); action = initComponent(ruleMd.getAction());
}
} }
private void initProcessor() throws Exception { private void initProcessor() throws Exception {
@ -131,9 +133,11 @@ class RuleActorMessageProcessor extends ComponentMsgProcessor<RuleId> {
} }
private void fetchPluginInfo() { private void fetchPluginInfo() {
PluginMetaData pluginMd = systemContext.getPluginService().findPluginByApiToken(ruleMd.getPluginToken()); if (!StringUtils.isEmpty(ruleMd.getPluginToken())) {
pluginTenantId = pluginMd.getTenantId(); PluginMetaData pluginMd = systemContext.getPluginService().findPluginByApiToken(ruleMd.getPluginToken());
pluginId = pluginMd.getId(); pluginTenantId = pluginMd.getTenantId();
pluginId = pluginMd.getId();
}
} }
protected void onRuleProcessingMsg(ActorContext context, RuleProcessingMsg msg) throws RuleException { protected void onRuleProcessingMsg(ActorContext context, RuleProcessingMsg msg) throws RuleException {
@ -162,25 +166,26 @@ class RuleActorMessageProcessor extends ComponentMsgProcessor<RuleId> {
inMsgMd = new RuleProcessingMetaData(); inMsgMd = new RuleProcessingMetaData();
} }
logger.debug("[{}] Going to convert in msg: {}", entityId, inMsg); logger.debug("[{}] Going to convert in msg: {}", entityId, inMsg);
Optional<RuleToPluginMsg<?>> ruleToPluginMsgOptional = action.convert(ruleCtx, inMsg, inMsgMd); if (action != null) {
if (ruleToPluginMsgOptional.isPresent()) { Optional<RuleToPluginMsg<?>> ruleToPluginMsgOptional = action.convert(ruleCtx, inMsg, inMsgMd);
RuleToPluginMsg<?> ruleToPluginMsg = ruleToPluginMsgOptional.get(); if (ruleToPluginMsgOptional.isPresent()) {
logger.debug("[{}] Device msg is converter to: {}", entityId, ruleToPluginMsg); RuleToPluginMsg<?> ruleToPluginMsg = ruleToPluginMsgOptional.get();
context.parent().tell(new RuleToPluginMsgWrapper(pluginTenantId, pluginId, tenantId, entityId, ruleToPluginMsg), context.self()); logger.debug("[{}] Device msg is converter to: {}", entityId, ruleToPluginMsg);
if (action.isOneWayAction()) { context.parent().tell(new RuleToPluginMsgWrapper(pluginTenantId, pluginId, tenantId, entityId, ruleToPluginMsg), context.self());
pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_TWO_WAY_ACTIONS); if (action.isOneWayAction()) {
} else { pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_TWO_WAY_ACTIONS);
pendingMsgMap.put(ruleToPluginMsg.getUid(), msg); } else {
scheduleMsgWithDelay(context, new RuleToPluginTimeoutMsg(ruleToPluginMsg.getUid()), systemContext.getPluginProcessingTimeout()); pendingMsgMap.put(ruleToPluginMsg.getUid(), msg);
scheduleMsgWithDelay(context, new RuleToPluginTimeoutMsg(ruleToPluginMsg.getUid()), systemContext.getPluginProcessingTimeout());
}
} }
} else { } else {
logger.debug("[{}] Nothing to send to plugin: {}", entityId, pluginId); logger.debug("[{}] Nothing to send to plugin: {}", entityId, pluginId);
pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_REQUEST_FROM_ACTIONS); pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_TWO_WAY_ACTIONS);
return;
} }
} }
public void onPluginMsg(ActorContext context, PluginToRuleMsg<?> msg) { void onPluginMsg(ActorContext context, PluginToRuleMsg<?> msg) {
RuleProcessingMsg pendingMsg = pendingMsgMap.remove(msg.getUid()); RuleProcessingMsg pendingMsg = pendingMsgMap.remove(msg.getUid());
if (pendingMsg != null) { if (pendingMsg != null) {
ChainProcessingContext ctx = pendingMsg.getCtx(); ChainProcessingContext ctx = pendingMsg.getCtx();
@ -196,7 +201,7 @@ class RuleActorMessageProcessor extends ComponentMsgProcessor<RuleId> {
} }
} }
public void onTimeoutMsg(ActorContext context, RuleToPluginTimeoutMsg msg) { void onTimeoutMsg(ActorContext context, RuleToPluginTimeoutMsg msg) {
RuleProcessingMsg pendingMsg = pendingMsgMap.remove(msg.getMsgId()); RuleProcessingMsg pendingMsg = pendingMsgMap.remove(msg.getMsgId());
if (pendingMsg != null) { if (pendingMsg != null) {
logger.debug("[{}] Processing timeout detected [{}]: {}", entityId, msg.getMsgId(), pendingMsg); logger.debug("[{}] Processing timeout detected [{}]: {}", entityId, msg.getMsgId(), pendingMsg);
@ -269,18 +274,16 @@ class RuleActorMessageProcessor extends ComponentMsgProcessor<RuleId> {
public void onActivate(ActorContext context) throws Exception { public void onActivate(ActorContext context) throws Exception {
logger.info("[{}] Going to process onActivate rule.", entityId); logger.info("[{}] Going to process onActivate rule.", entityId);
this.state = ComponentLifecycleState.ACTIVE; this.state = ComponentLifecycleState.ACTIVE;
if (action != null) { if (filters != null) {
if (filters != null) { filters.forEach(RuleLifecycleComponent::resume);
filters.forEach(f -> f.resume());
} else {
initFilters();
}
if (processor != null) { if (processor != null) {
processor.resume(); processor.resume();
} else { } else {
initProcessor(); initProcessor();
} }
action.resume(); if (action != null) {
action.resume();
}
logger.info("[{}] Rule resumed.", entityId); logger.info("[{}] Rule resumed.", entityId);
} else { } else {
start(); start();

7
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleService.java

@ -91,7 +91,9 @@ public class BaseRuleService extends AbstractEntityService implements RuleServic
if (rule.getProcessor() != null && !rule.getProcessor().isNull()) { if (rule.getProcessor() != null && !rule.getProcessor().isNull()) {
validateComponentJson(rule.getProcessor(), ComponentType.PROCESSOR); validateComponentJson(rule.getProcessor(), ComponentType.PROCESSOR);
} }
validateComponentJson(rule.getAction(), ComponentType.ACTION); if (rule.getAction() != null && !rule.getAction().isNull()) {
validateComponentJson(rule.getAction(), ComponentType.ACTION);
}
validateRuleAndPluginState(rule); validateRuleAndPluginState(rule);
return ruleDao.save(rule); return ruleDao.save(rule);
} }
@ -129,6 +131,9 @@ public class BaseRuleService extends AbstractEntityService implements RuleServic
} }
private void validateRuleAndPluginState(RuleMetaData rule) { private void validateRuleAndPluginState(RuleMetaData rule) {
if (org.springframework.util.StringUtils.isEmpty(rule.getPluginToken())) {
return;
}
PluginMetaData pluginMd = pluginService.findPluginByApiToken(rule.getPluginToken()); PluginMetaData pluginMd = pluginService.findPluginByApiToken(rule.getPluginToken());
if (pluginMd == null) { if (pluginMd == null) {
throw new IncorrectParameterException("Rule points to non-existent plugin!"); throw new IncorrectParameterException("Rule points to non-existent plugin!");

Loading…
Cancel
Save