96 changed files with 1238 additions and 1558 deletions
@ -1,117 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import akka.actor.ActorRef; |
|
||||
import org.thingsboard.server.common.msg.core.RuleEngineError; |
|
||||
import org.thingsboard.server.common.msg.core.RuleEngineErrorMsg; |
|
||||
import org.thingsboard.server.common.msg.device.ToDeviceActorMsg; |
|
||||
import org.thingsboard.server.common.msg.session.ToDeviceMsg; |
|
||||
import org.thingsboard.server.extensions.api.device.DeviceAttributes; |
|
||||
import org.thingsboard.server.extensions.api.device.DeviceMetaData; |
|
||||
|
|
||||
public class ChainProcessingContext { |
|
||||
|
|
||||
private final ChainProcessingMetaData md; |
|
||||
private final int index; |
|
||||
private final RuleEngineError error; |
|
||||
private ToDeviceMsg response; |
|
||||
|
|
||||
|
|
||||
public ChainProcessingContext(ChainProcessingMetaData md) { |
|
||||
super(); |
|
||||
this.md = md; |
|
||||
this.index = 0; |
|
||||
this.error = RuleEngineError.NO_RULES; |
|
||||
} |
|
||||
|
|
||||
private ChainProcessingContext(ChainProcessingContext other, int indexOffset, RuleEngineError error) { |
|
||||
super(); |
|
||||
this.md = other.md; |
|
||||
this.index = other.index + indexOffset; |
|
||||
this.error = error; |
|
||||
this.response = other.response; |
|
||||
|
|
||||
if (this.index < 0 || this.index >= this.md.chain.size()) { |
|
||||
throw new IllegalArgumentException("Can't apply offset " + indexOffset + " to the chain!"); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public ActorRef getDeviceActor() { |
|
||||
return md.originator; |
|
||||
} |
|
||||
|
|
||||
public ActorRef getCurrentActor() { |
|
||||
return md.chain.getRuleActorMd(index).getActorRef(); |
|
||||
} |
|
||||
|
|
||||
public boolean hasNext() { |
|
||||
return (getChainLength() - 1) > index; |
|
||||
} |
|
||||
|
|
||||
public boolean isFailure() { |
|
||||
return (error != null && error.isCritical()) || (response != null && !response.isSuccess()); |
|
||||
} |
|
||||
|
|
||||
public ChainProcessingContext getNext() { |
|
||||
return new ChainProcessingContext(this, 1, this.error); |
|
||||
} |
|
||||
|
|
||||
public ChainProcessingContext withError(RuleEngineError error) { |
|
||||
if (error != null && (this.error == null || this.error.getPriority() < error.getPriority())) { |
|
||||
return new ChainProcessingContext(this, 0, error); |
|
||||
} else { |
|
||||
return this; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public int getChainLength() { |
|
||||
return md.chain.size(); |
|
||||
} |
|
||||
|
|
||||
public ToDeviceActorMsg getInMsg() { |
|
||||
return md.inMsg; |
|
||||
} |
|
||||
|
|
||||
public DeviceMetaData getDeviceMetaData() { |
|
||||
return md.deviceMetaData; |
|
||||
} |
|
||||
|
|
||||
public String getDeviceName() { |
|
||||
return md.deviceMetaData.getDeviceName(); |
|
||||
} |
|
||||
|
|
||||
public String getDeviceType() { |
|
||||
return md.deviceMetaData.getDeviceType(); |
|
||||
} |
|
||||
|
|
||||
public DeviceAttributes getAttributes() { |
|
||||
return md.deviceMetaData.getDeviceAttributes(); |
|
||||
} |
|
||||
|
|
||||
public ToDeviceMsg getResponse() { |
|
||||
return response; |
|
||||
} |
|
||||
|
|
||||
public void mergeResponse(ToDeviceMsg response) { |
|
||||
// TODO add merge logic
|
|
||||
this.response = response; |
|
||||
} |
|
||||
|
|
||||
public RuleEngineErrorMsg getError() { |
|
||||
return new RuleEngineErrorMsg(md.inMsg.getPayload().getMsgType(), error); |
|
||||
} |
|
||||
} |
|
||||
@ -1,41 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import akka.actor.ActorRef; |
|
||||
import org.thingsboard.server.common.msg.device.ToDeviceActorMsg; |
|
||||
import org.thingsboard.server.extensions.api.device.DeviceMetaData; |
|
||||
|
|
||||
/** |
|
||||
* Immutable part of chain processing data; |
|
||||
* |
|
||||
* @author ashvayka |
|
||||
*/ |
|
||||
public final class ChainProcessingMetaData { |
|
||||
|
|
||||
final RuleActorChain chain; |
|
||||
final ToDeviceActorMsg inMsg; |
|
||||
final ActorRef originator; |
|
||||
final DeviceMetaData deviceMetaData; |
|
||||
|
|
||||
public ChainProcessingMetaData(RuleActorChain chain, ToDeviceActorMsg inMsg, DeviceMetaData deviceMetaData, ActorRef originator) { |
|
||||
super(); |
|
||||
this.chain = chain; |
|
||||
this.inMsg = inMsg; |
|
||||
this.originator = originator; |
|
||||
this.deviceMetaData = deviceMetaData; |
|
||||
} |
|
||||
} |
|
||||
@ -1,43 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
public class ComplexRuleActorChain implements RuleActorChain { |
|
||||
|
|
||||
private final RuleActorChain systemChain; |
|
||||
private final RuleActorChain tenantChain; |
|
||||
|
|
||||
public ComplexRuleActorChain(RuleActorChain systemChain, RuleActorChain tenantChain) { |
|
||||
super(); |
|
||||
this.systemChain = systemChain; |
|
||||
this.tenantChain = tenantChain; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public int size() { |
|
||||
return systemChain.size() + tenantChain.size(); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public RuleActorMetaData getRuleActorMd(int index) { |
|
||||
if (index < systemChain.size()) { |
|
||||
return systemChain.getRuleActorMd(index); |
|
||||
} else { |
|
||||
return tenantChain.getRuleActorMd(index - systemChain.size()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,20 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
public class CompoundRuleActorChain { |
|
||||
|
|
||||
} |
|
||||
@ -1,90 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.service.ComponentActor; |
|
||||
import org.thingsboard.server.actors.service.ContextBasedCreator; |
|
||||
import org.thingsboard.server.actors.stats.StatsPersistTick; |
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; |
|
||||
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.PluginToRuleMsg; |
|
||||
|
|
||||
public class RuleActor extends ComponentActor<RuleId, RuleActorMessageProcessor> { |
|
||||
|
|
||||
private RuleActor(ActorSystemContext systemContext, TenantId tenantId, RuleId ruleId) { |
|
||||
super(systemContext, tenantId, ruleId); |
|
||||
setProcessor(new RuleActorMessageProcessor(tenantId, ruleId, systemContext, logger)); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onReceive(Object msg) throws Exception { |
|
||||
logger.debug("[{}] Received message: {}", id, msg); |
|
||||
if (msg instanceof RuleProcessingMsg) { |
|
||||
try { |
|
||||
processor.onRuleProcessingMsg(context(), (RuleProcessingMsg) msg); |
|
||||
increaseMessagesProcessedCount(); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onDeviceMsg", e); |
|
||||
} |
|
||||
} else if (msg instanceof PluginToRuleMsg<?>) { |
|
||||
try { |
|
||||
processor.onPluginMsg(context(), (PluginToRuleMsg<?>) msg); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onPluginMsg", e); |
|
||||
} |
|
||||
} else if (msg instanceof ComponentLifecycleMsg) { |
|
||||
onComponentLifecycleMsg((ComponentLifecycleMsg) msg); |
|
||||
} else if (msg instanceof ClusterEventMsg) { |
|
||||
onClusterEventMsg((ClusterEventMsg) msg); |
|
||||
} else if (msg instanceof RuleToPluginTimeoutMsg) { |
|
||||
try { |
|
||||
processor.onTimeoutMsg(context(), (RuleToPluginTimeoutMsg) msg); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onTimeoutMsg", e); |
|
||||
} |
|
||||
} else if (msg instanceof StatsPersistTick) { |
|
||||
onStatsPersistTick(id); |
|
||||
} else { |
|
||||
logger.debug("[{}][{}] Unknown msg type.", tenantId, id, msg.getClass().getName()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public static class ActorCreator extends ContextBasedCreator<RuleActor> { |
|
||||
private static final long serialVersionUID = 1L; |
|
||||
|
|
||||
private final TenantId tenantId; |
|
||||
private final RuleId ruleId; |
|
||||
|
|
||||
public ActorCreator(ActorSystemContext context, TenantId tenantId, RuleId ruleId) { |
|
||||
super(context); |
|
||||
this.tenantId = tenantId; |
|
||||
this.ruleId = ruleId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public RuleActor create() throws Exception { |
|
||||
return new RuleActor(context, tenantId, ruleId); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected long getErrorPersistFrequency() { |
|
||||
return systemContext.getRuleErrorPersistFrequency(); |
|
||||
} |
|
||||
} |
|
||||
@ -1,345 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import java.util.*; |
|
||||
|
|
||||
import com.fasterxml.jackson.core.JsonProcessingException; |
|
||||
import org.springframework.util.StringUtils; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.plugin.RuleToPluginMsgWrapper; |
|
||||
import org.thingsboard.server.actors.shared.ComponentMsgProcessor; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleState; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
import org.thingsboard.server.common.data.rule.RuleMetaData; |
|
||||
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; |
|
||||
import org.thingsboard.server.common.msg.core.BasicRequest; |
|
||||
import org.thingsboard.server.common.msg.core.BasicStatusCodeResponse; |
|
||||
import org.thingsboard.server.common.msg.core.RuleEngineError; |
|
||||
import org.thingsboard.server.common.msg.device.ToDeviceActorMsg; |
|
||||
import org.thingsboard.server.common.msg.session.MsgType; |
|
||||
import org.thingsboard.server.common.msg.session.ToDeviceMsg; |
|
||||
import org.thingsboard.server.common.msg.session.ex.ProcessingTimeoutException; |
|
||||
import org.thingsboard.server.extensions.api.rules.*; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginAction; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.PluginToRuleMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.RuleToPluginMsg; |
|
||||
|
|
||||
import com.fasterxml.jackson.databind.JsonNode; |
|
||||
|
|
||||
import akka.actor.ActorContext; |
|
||||
import akka.actor.ActorRef; |
|
||||
import akka.event.LoggingAdapter; |
|
||||
|
|
||||
class RuleActorMessageProcessor extends ComponentMsgProcessor<RuleId> { |
|
||||
|
|
||||
private final RuleProcessingContext ruleCtx; |
|
||||
private final Map<UUID, RuleProcessingMsg> pendingMsgMap; |
|
||||
|
|
||||
private RuleMetaData ruleMd; |
|
||||
private ComponentLifecycleState state; |
|
||||
private List<RuleFilter> filters; |
|
||||
private RuleProcessor processor; |
|
||||
private PluginAction action; |
|
||||
|
|
||||
private TenantId pluginTenantId; |
|
||||
private PluginId pluginId; |
|
||||
|
|
||||
protected RuleActorMessageProcessor(TenantId tenantId, RuleId ruleId, ActorSystemContext systemContext, LoggingAdapter logger) { |
|
||||
super(systemContext, logger, tenantId, ruleId); |
|
||||
this.pendingMsgMap = new HashMap<>(); |
|
||||
this.ruleCtx = new RuleProcessingContext(systemContext, ruleId); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void start(ActorContext context) throws Exception { |
|
||||
logger.info("[{}][{}] Starting rule actor.", entityId, tenantId); |
|
||||
ruleMd = systemContext.getRuleService().findRuleById(entityId); |
|
||||
if (ruleMd == null) { |
|
||||
throw new RuleInitializationException("Rule not found!"); |
|
||||
} |
|
||||
state = ruleMd.getState(); |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
logger.info("[{}] Rule is active. Going to initialize rule components.", entityId); |
|
||||
initComponent(); |
|
||||
} else { |
|
||||
logger.info("[{}] Rule is suspended. Skipping rule components initialization.", entityId); |
|
||||
} |
|
||||
|
|
||||
logger.info("[{}][{}] Started rule actor.", entityId, tenantId); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void stop(ActorContext context) throws Exception { |
|
||||
onStop(); |
|
||||
} |
|
||||
|
|
||||
|
|
||||
private void initComponent() throws RuleException { |
|
||||
try { |
|
||||
if (!ruleMd.getFilters().isArray()) { |
|
||||
throw new RuntimeException("Filters are not array!"); |
|
||||
} |
|
||||
fetchPluginInfo(); |
|
||||
initFilters(); |
|
||||
initProcessor(); |
|
||||
initAction(); |
|
||||
} catch (RuntimeException e) { |
|
||||
throw new RuleInitializationException("Unknown runtime exception!", e); |
|
||||
} catch (InstantiationException e) { |
|
||||
throw new RuleInitializationException("No default constructor for rule implementation!", e); |
|
||||
} catch (IllegalAccessException e) { |
|
||||
throw new RuleInitializationException("Illegal Access Exception during rule initialization!", e); |
|
||||
} catch (ClassNotFoundException e) { |
|
||||
throw new RuleInitializationException("Rule Class not found!", e); |
|
||||
} catch (Exception e) { |
|
||||
throw new RuleException(e.getMessage(), e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void initAction() throws Exception { |
|
||||
if (ruleMd.getAction() != null && !ruleMd.getAction().isNull()) { |
|
||||
action = initComponent(ruleMd.getAction()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void initProcessor() throws Exception { |
|
||||
if (ruleMd.getProcessor() != null && !ruleMd.getProcessor().isNull()) { |
|
||||
processor = initComponent(ruleMd.getProcessor()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void initFilters() throws Exception { |
|
||||
filters = new ArrayList<>(ruleMd.getFilters().size()); |
|
||||
for (int i = 0; i < ruleMd.getFilters().size(); i++) { |
|
||||
filters.add(initComponent(ruleMd.getFilters().get(i))); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void fetchPluginInfo() { |
|
||||
if (!StringUtils.isEmpty(ruleMd.getPluginToken())) { |
|
||||
PluginMetaData pluginMd = systemContext.getPluginService().findPluginByApiToken(ruleMd.getPluginToken()); |
|
||||
pluginTenantId = pluginMd.getTenantId(); |
|
||||
pluginId = pluginMd.getId(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
protected void onRuleProcessingMsg(ActorContext context, RuleProcessingMsg msg) throws RuleException { |
|
||||
if (state != ComponentLifecycleState.ACTIVE) { |
|
||||
pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_ACTIVE_RULES); |
|
||||
return; |
|
||||
} |
|
||||
ChainProcessingContext chainCtx = msg.getCtx(); |
|
||||
ToDeviceActorMsg inMsg = chainCtx.getInMsg(); |
|
||||
|
|
||||
ruleCtx.update(inMsg, chainCtx.getDeviceMetaData()); |
|
||||
|
|
||||
logger.debug("[{}] Going to filter in msg: {}", entityId, inMsg); |
|
||||
for (RuleFilter filter : filters) { |
|
||||
if (!filter.filter(ruleCtx, inMsg)) { |
|
||||
logger.debug("[{}] In msg is NOT valid for processing by current rule: {}", entityId, inMsg); |
|
||||
pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_FILTERS_MATCHED); |
|
||||
return; |
|
||||
} |
|
||||
} |
|
||||
RuleProcessingMetaData inMsgMd; |
|
||||
if (processor != null) { |
|
||||
logger.debug("[{}] Going to process in msg: {}", entityId, inMsg); |
|
||||
inMsgMd = processor.process(ruleCtx, inMsg); |
|
||||
} else { |
|
||||
inMsgMd = new RuleProcessingMetaData(); |
|
||||
} |
|
||||
logger.debug("[{}] Going to convert in msg: {}", entityId, inMsg); |
|
||||
if (action != null) { |
|
||||
Optional<RuleToPluginMsg<?>> ruleToPluginMsgOptional = action.convert(ruleCtx, inMsg, inMsgMd); |
|
||||
if (ruleToPluginMsgOptional.isPresent()) { |
|
||||
RuleToPluginMsg<?> ruleToPluginMsg = ruleToPluginMsgOptional.get(); |
|
||||
logger.debug("[{}] Device msg is converted to: {}", entityId, ruleToPluginMsg); |
|
||||
context.parent().tell(new RuleToPluginMsgWrapper(pluginTenantId, pluginId, tenantId, entityId, ruleToPluginMsg), context.self()); |
|
||||
if (action.isOneWayAction()) { |
|
||||
pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_TWO_WAY_ACTIONS); |
|
||||
return; |
|
||||
} else { |
|
||||
pendingMsgMap.put(ruleToPluginMsg.getUid(), msg); |
|
||||
scheduleMsgWithDelay(context, new RuleToPluginTimeoutMsg(ruleToPluginMsg.getUid()), systemContext.getPluginProcessingTimeout()); |
|
||||
return; |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
logger.debug("[{}] Nothing to send to plugin: {}", entityId, pluginId); |
|
||||
pushToNextRule(context, msg.getCtx(), RuleEngineError.NO_TWO_WAY_ACTIONS); |
|
||||
} |
|
||||
|
|
||||
void onPluginMsg(ActorContext context, PluginToRuleMsg<?> msg) { |
|
||||
RuleProcessingMsg pendingMsg = pendingMsgMap.remove(msg.getUid()); |
|
||||
if (pendingMsg != null) { |
|
||||
ChainProcessingContext ctx = pendingMsg.getCtx(); |
|
||||
Optional<ToDeviceMsg> ruleResponseOptional = action.convert(msg); |
|
||||
if (ruleResponseOptional.isPresent()) { |
|
||||
ctx.mergeResponse(ruleResponseOptional.get()); |
|
||||
pushToNextRule(context, ctx, null); |
|
||||
} else { |
|
||||
pushToNextRule(context, ctx, RuleEngineError.NO_RESPONSE_FROM_ACTIONS); |
|
||||
} |
|
||||
} else { |
|
||||
logger.warning("[{}] Processing timeout detected: [{}]", entityId, msg.getUid()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
void onTimeoutMsg(ActorContext context, RuleToPluginTimeoutMsg msg) { |
|
||||
RuleProcessingMsg pendingMsg = pendingMsgMap.remove(msg.getMsgId()); |
|
||||
if (pendingMsg != null) { |
|
||||
logger.debug("[{}] Processing timeout detected [{}]: {}", entityId, msg.getMsgId(), pendingMsg); |
|
||||
ChainProcessingContext ctx = pendingMsg.getCtx(); |
|
||||
pushToNextRule(context, ctx, RuleEngineError.PLUGIN_TIMEOUT); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void pushToNextRule(ActorContext context, ChainProcessingContext ctx, RuleEngineError error) { |
|
||||
if (error != null) { |
|
||||
ctx = ctx.withError(error); |
|
||||
} |
|
||||
if (ctx.isFailure()) { |
|
||||
logger.debug("[{}][{}] Forwarding processing chain to device actor due to failure.", ruleMd.getId(), ctx.getInMsg().getDeviceId()); |
|
||||
ctx.getDeviceActor().tell(new RulesProcessedMsg(ctx), ActorRef.noSender()); |
|
||||
} else if (!ctx.hasNext()) { |
|
||||
logger.debug("[{}][{}] Forwarding processing chain to device actor due to end of chain.", ruleMd.getId(), ctx.getInMsg().getDeviceId()); |
|
||||
ctx.getDeviceActor().tell(new RulesProcessedMsg(ctx), ActorRef.noSender()); |
|
||||
} else { |
|
||||
logger.debug("[{}][{}] Forwarding processing chain to next rule actor.", ruleMd.getId(), ctx.getInMsg().getDeviceId()); |
|
||||
ChainProcessingContext nextTask = ctx.getNext(); |
|
||||
nextTask.getCurrentActor().tell(new RuleProcessingMsg(nextTask), context.self()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onCreated(ActorContext context) { |
|
||||
logger.info("[{}] Going to process onCreated rule.", entityId); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onUpdate(ActorContext context) throws RuleException { |
|
||||
RuleMetaData oldRuleMd = ruleMd; |
|
||||
ruleMd = systemContext.getRuleService().findRuleById(entityId); |
|
||||
logger.info("[{}] Rule configuration was updated from {} to {}.", entityId, oldRuleMd, ruleMd); |
|
||||
try { |
|
||||
fetchPluginInfo(); |
|
||||
if (filters == null || !Objects.equals(oldRuleMd.getFilters(), ruleMd.getFilters())) { |
|
||||
logger.info("[{}] Rule filters require restart due to json change from {} to {}.", |
|
||||
entityId, mapper.writeValueAsString(oldRuleMd.getFilters()), mapper.writeValueAsString(ruleMd.getFilters())); |
|
||||
stopFilters(); |
|
||||
initFilters(); |
|
||||
} |
|
||||
if (processor == null || !Objects.equals(oldRuleMd.getProcessor(), ruleMd.getProcessor())) { |
|
||||
logger.info("[{}] Rule processor require restart due to configuration change.", entityId); |
|
||||
stopProcessor(); |
|
||||
initProcessor(); |
|
||||
} |
|
||||
if (action == null || !Objects.equals(oldRuleMd.getAction(), ruleMd.getAction())) { |
|
||||
logger.info("[{}] Rule action require restart due to configuration change.", entityId); |
|
||||
stopAction(); |
|
||||
initAction(); |
|
||||
} |
|
||||
} catch (RuntimeException e) { |
|
||||
throw new RuleInitializationException("Unknown runtime exception!", e); |
|
||||
} catch (InstantiationException e) { |
|
||||
throw new RuleInitializationException("No default constructor for rule implementation!", e); |
|
||||
} catch (IllegalAccessException e) { |
|
||||
throw new RuleInitializationException("Illegal Access Exception during rule initialization!", e); |
|
||||
} catch (ClassNotFoundException e) { |
|
||||
throw new RuleInitializationException("Rule Class not found!", e); |
|
||||
} catch (JsonProcessingException e) { |
|
||||
throw new RuleInitializationException("Rule configuration is invalid!", e); |
|
||||
} catch (Exception e) { |
|
||||
throw new RuleInitializationException(e.getMessage(), e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onActivate(ActorContext context) throws Exception { |
|
||||
logger.info("[{}] Going to process onActivate rule.", entityId); |
|
||||
this.state = ComponentLifecycleState.ACTIVE; |
|
||||
if (filters != null) { |
|
||||
filters.forEach(RuleLifecycleComponent::resume); |
|
||||
if (processor != null) { |
|
||||
processor.resume(); |
|
||||
} else { |
|
||||
initProcessor(); |
|
||||
} |
|
||||
if (action != null) { |
|
||||
action.resume(); |
|
||||
} |
|
||||
logger.info("[{}] Rule resumed.", entityId); |
|
||||
} else { |
|
||||
start(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onSuspend(ActorContext context) { |
|
||||
logger.info("[{}] Going to process onSuspend rule.", entityId); |
|
||||
this.state = ComponentLifecycleState.SUSPENDED; |
|
||||
if (filters != null) { |
|
||||
filters.forEach(f -> f.suspend()); |
|
||||
} |
|
||||
if (processor != null) { |
|
||||
processor.suspend(); |
|
||||
} |
|
||||
if (action != null) { |
|
||||
action.suspend(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onStop(ActorContext context) { |
|
||||
logger.info("[{}] Going to process onStop rule.", entityId); |
|
||||
onStop(); |
|
||||
scheduleMsgWithDelay(context, new RuleTerminationMsg(entityId), systemContext.getRuleActorTerminationDelay()); |
|
||||
} |
|
||||
|
|
||||
private void onStop() { |
|
||||
this.state = ComponentLifecycleState.SUSPENDED; |
|
||||
stopFilters(); |
|
||||
stopProcessor(); |
|
||||
stopAction(); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onClusterEventMsg(ClusterEventMsg msg) throws Exception { |
|
||||
//Do nothing
|
|
||||
} |
|
||||
|
|
||||
private void stopAction() { |
|
||||
if (action != null) { |
|
||||
action.stop(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void stopProcessor() { |
|
||||
if (processor != null) { |
|
||||
processor.stop(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void stopFilters() { |
|
||||
if (filters != null) { |
|
||||
filters.forEach(f -> f.stop()); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,107 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import java.util.Comparator; |
|
||||
|
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
|
|
||||
import akka.actor.ActorRef; |
|
||||
|
|
||||
public class RuleActorMetaData { |
|
||||
|
|
||||
private final RuleId ruleId; |
|
||||
private final boolean systemRule; |
|
||||
private final int weight; |
|
||||
private final ActorRef actorRef; |
|
||||
|
|
||||
public static final Comparator<RuleActorMetaData> RULE_ACTOR_MD_COMPARATOR = new Comparator<RuleActorMetaData>() { |
|
||||
|
|
||||
@Override |
|
||||
public int compare(RuleActorMetaData r1, RuleActorMetaData r2) { |
|
||||
if (r1.isSystemRule() && !r2.isSystemRule()) { |
|
||||
return 1; |
|
||||
} else if (!r1.isSystemRule() && r2.isSystemRule()) { |
|
||||
return -1; |
|
||||
} else { |
|
||||
return Integer.compare(r2.getWeight(), r1.getWeight()); |
|
||||
} |
|
||||
} |
|
||||
}; |
|
||||
|
|
||||
public static RuleActorMetaData systemRule(RuleId ruleId, int weight, ActorRef actorRef) { |
|
||||
return new RuleActorMetaData(ruleId, true, weight, actorRef); |
|
||||
} |
|
||||
|
|
||||
public static RuleActorMetaData tenantRule(RuleId ruleId, int weight, ActorRef actorRef) { |
|
||||
return new RuleActorMetaData(ruleId, false, weight, actorRef); |
|
||||
} |
|
||||
|
|
||||
private RuleActorMetaData(RuleId ruleId, boolean systemRule, int weight, ActorRef actorRef) { |
|
||||
super(); |
|
||||
this.ruleId = ruleId; |
|
||||
this.systemRule = systemRule; |
|
||||
this.weight = weight; |
|
||||
this.actorRef = actorRef; |
|
||||
} |
|
||||
|
|
||||
public RuleId getRuleId() { |
|
||||
return ruleId; |
|
||||
} |
|
||||
|
|
||||
public boolean isSystemRule() { |
|
||||
return systemRule; |
|
||||
} |
|
||||
|
|
||||
public int getWeight() { |
|
||||
return weight; |
|
||||
} |
|
||||
|
|
||||
public ActorRef getActorRef() { |
|
||||
return actorRef; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public int hashCode() { |
|
||||
final int prime = 31; |
|
||||
int result = 1; |
|
||||
result = prime * result + ((ruleId == null) ? 0 : ruleId.hashCode()); |
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public boolean equals(Object obj) { |
|
||||
if (this == obj) |
|
||||
return true; |
|
||||
if (obj == null) |
|
||||
return false; |
|
||||
if (getClass() != obj.getClass()) |
|
||||
return false; |
|
||||
RuleActorMetaData other = (RuleActorMetaData) obj; |
|
||||
if (ruleId == null) { |
|
||||
if (other.ruleId != null) |
|
||||
return false; |
|
||||
} else if (!ruleId.equals(other.ruleId)) |
|
||||
return false; |
|
||||
return true; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public String toString() { |
|
||||
return "RuleActorMetaData [ruleId=" + ruleId + ", systemRule=" + systemRule + ", weight=" + weight + ", actorRef=" + actorRef + "]"; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,33 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor; |
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
|
|
||||
import akka.event.LoggingAdapter; |
|
||||
|
|
||||
public class RuleContextAwareMsgProcessor extends AbstractContextAwareMsgProcessor { |
|
||||
|
|
||||
private final RuleId ruleId; |
|
||||
|
|
||||
protected RuleContextAwareMsgProcessor(ActorSystemContext systemContext, LoggingAdapter logger, RuleId ruleId) { |
|
||||
super(systemContext, logger); |
|
||||
this.ruleId = ruleId; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,115 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.common.data.Event; |
|
||||
import org.thingsboard.server.common.data.alarm.Alarm; |
|
||||
import org.thingsboard.server.common.data.alarm.AlarmId; |
|
||||
import org.thingsboard.server.common.data.id.*; |
|
||||
import org.thingsboard.server.dao.alarm.AlarmService; |
|
||||
import org.thingsboard.server.dao.event.EventService; |
|
||||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|
||||
import org.thingsboard.server.common.msg.device.ToDeviceActorMsg; |
|
||||
import org.thingsboard.server.extensions.api.device.DeviceMetaData; |
|
||||
import org.thingsboard.server.extensions.api.rules.RuleContext; |
|
||||
|
|
||||
import java.util.Optional; |
|
||||
import java.util.concurrent.ExecutionException; |
|
||||
|
|
||||
public class RuleProcessingContext implements RuleContext { |
|
||||
|
|
||||
private final TimeseriesService tsService; |
|
||||
private final EventService eventService; |
|
||||
private final AlarmService alarmService; |
|
||||
private final RuleId ruleId; |
|
||||
private TenantId tenantId; |
|
||||
private CustomerId customerId; |
|
||||
private DeviceId deviceId; |
|
||||
private DeviceMetaData deviceMetaData; |
|
||||
|
|
||||
RuleProcessingContext(ActorSystemContext systemContext, RuleId ruleId) { |
|
||||
this.tsService = systemContext.getTsService(); |
|
||||
this.eventService = systemContext.getEventService(); |
|
||||
this.alarmService = systemContext.getAlarmService(); |
|
||||
this.ruleId = ruleId; |
|
||||
} |
|
||||
|
|
||||
void update(ToDeviceActorMsg toDeviceActorMsg, DeviceMetaData deviceMetaData) { |
|
||||
this.tenantId = toDeviceActorMsg.getTenantId(); |
|
||||
this.customerId = toDeviceActorMsg.getCustomerId(); |
|
||||
this.deviceId = toDeviceActorMsg.getDeviceId(); |
|
||||
this.deviceMetaData = deviceMetaData; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public RuleId getRuleId() { |
|
||||
return ruleId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public DeviceMetaData getDeviceMetaData() { |
|
||||
return deviceMetaData; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Event save(Event event) { |
|
||||
checkEvent(event); |
|
||||
return eventService.save(event); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Optional<Event> saveIfNotExists(Event event) { |
|
||||
checkEvent(event); |
|
||||
return eventService.saveIfNotExists(event); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Optional<Event> findEvent(String eventType, String eventUid) { |
|
||||
return eventService.findEvent(tenantId, deviceId, eventType, eventUid); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Alarm createOrUpdateAlarm(Alarm alarm) { |
|
||||
alarm.setTenantId(tenantId); |
|
||||
return alarmService.createOrUpdateAlarm(alarm); |
|
||||
} |
|
||||
|
|
||||
public Optional<Alarm> findLatestAlarm(EntityId originator, String alarmType) { |
|
||||
try { |
|
||||
return Optional.ofNullable(alarmService.findLatestByOriginatorAndType(tenantId, originator, alarmType).get()); |
|
||||
} catch (InterruptedException | ExecutionException e) { |
|
||||
throw new RuntimeException("Failed to lookup alarm!", e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Boolean> clearAlarm(AlarmId alarmId, long clearTs) { |
|
||||
return alarmService.clearAlarm(alarmId, clearTs); |
|
||||
} |
|
||||
|
|
||||
private void checkEvent(Event event) { |
|
||||
if (event.getTenantId() == null) { |
|
||||
event.setTenantId(tenantId); |
|
||||
} else if (!tenantId.equals(event.getTenantId())) { |
|
||||
throw new IllegalArgumentException("Invalid Tenant id!"); |
|
||||
} |
|
||||
if (event.getEntityId() == null) { |
|
||||
event.setEntityId(deviceId); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,31 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
public class RuleProcessingMsg { |
|
||||
|
|
||||
private final ChainProcessingContext ctx; |
|
||||
|
|
||||
public RuleProcessingMsg(ChainProcessingContext ctx) { |
|
||||
super(); |
|
||||
this.ctx = ctx; |
|
||||
} |
|
||||
|
|
||||
public ChainProcessingContext getCtx() { |
|
||||
return ctx; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,30 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import org.thingsboard.server.actors.shared.ActorTerminationMsg; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
|
|
||||
/** |
|
||||
* @author Andrew Shvayka |
|
||||
*/ |
|
||||
public class RuleTerminationMsg extends ActorTerminationMsg<RuleId> { |
|
||||
|
|
||||
public RuleTerminationMsg(RuleId id) { |
|
||||
super(id); |
|
||||
} |
|
||||
} |
|
||||
@ -1,36 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import java.io.Serializable; |
|
||||
import java.util.UUID; |
|
||||
|
|
||||
public class RuleToPluginTimeoutMsg implements Serializable { |
|
||||
|
|
||||
private static final long serialVersionUID = 1L; |
|
||||
|
|
||||
private final UUID msgId; |
|
||||
|
|
||||
public RuleToPluginTimeoutMsg(UUID msgId) { |
|
||||
super(); |
|
||||
this.msgId = msgId; |
|
||||
} |
|
||||
|
|
||||
public UUID getMsgId() { |
|
||||
return msgId; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,39 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.rule; |
|
||||
|
|
||||
import java.util.ArrayList; |
|
||||
import java.util.List; |
|
||||
import java.util.Set; |
|
||||
|
|
||||
public class SimpleRuleActorChain implements RuleActorChain { |
|
||||
|
|
||||
private final List<RuleActorMetaData> rules; |
|
||||
|
|
||||
public SimpleRuleActorChain(Set<RuleActorMetaData> ruleSet) { |
|
||||
rules = new ArrayList<>(ruleSet); |
|
||||
rules.sort(RuleActorMetaData.RULE_ACTOR_MD_COMPARATOR); |
|
||||
} |
|
||||
|
|
||||
public int size() { |
|
||||
return rules.size(); |
|
||||
} |
|
||||
|
|
||||
public RuleActorMetaData getRuleActorMd(int index) { |
|
||||
return rules.get(index); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,15 +1,33 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
package org.thingsboard.server.actors.ruleChain; |
package org.thingsboard.server.actors.ruleChain; |
||||
|
|
||||
import akka.actor.ActorRef; |
import akka.actor.ActorRef; |
||||
import lombok.Data; |
import lombok.Data; |
||||
import org.thingsboard.server.common.data.id.RuleNodeId; |
import org.thingsboard.server.common.data.id.RuleNodeId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.rule.RuleNode; |
||||
|
|
||||
/** |
/** |
||||
* Created by ashvayka on 19.03.18. |
* Created by ashvayka on 19.03.18. |
||||
*/ |
*/ |
||||
@Data |
@Data |
||||
final class RuleNodeCtx { |
final class RuleNodeCtx { |
||||
|
private final TenantId tenantId; |
||||
private final ActorRef chainActor; |
private final ActorRef chainActor; |
||||
private final ActorRef self; |
private final ActorRef selfActor; |
||||
private final RuleNodeId selfId; |
private final RuleNode self; |
||||
} |
} |
||||
|
|||||
@ -1,135 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.shared.rule; |
|
||||
|
|
||||
import akka.actor.ActorContext; |
|
||||
import akka.actor.ActorRef; |
|
||||
import akka.actor.Props; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.rule.RuleActor; |
|
||||
import org.thingsboard.server.actors.rule.RuleActorChain; |
|
||||
import org.thingsboard.server.actors.rule.RuleActorMetaData; |
|
||||
import org.thingsboard.server.actors.rule.SimpleRuleActorChain; |
|
||||
import org.thingsboard.server.actors.service.ContextAwareActor; |
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.PageDataIterable; |
|
||||
import org.thingsboard.server.common.data.page.PageDataIterable.FetchFunction; |
|
||||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
|
||||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleState; |
|
||||
import org.thingsboard.server.common.data.rule.RuleMetaData; |
|
||||
import org.thingsboard.server.dao.rule.RuleService; |
|
||||
|
|
||||
import java.util.*; |
|
||||
|
|
||||
@Slf4j |
|
||||
public abstract class RuleManager { |
|
||||
|
|
||||
protected final ActorSystemContext systemContext; |
|
||||
protected final RuleService ruleService; |
|
||||
protected final Map<RuleId, ActorRef> ruleActors; |
|
||||
protected final TenantId tenantId; |
|
||||
|
|
||||
private Map<RuleMetaData, RuleActorMetaData> ruleMap; |
|
||||
private RuleActorChain ruleChain; |
|
||||
|
|
||||
public RuleManager(ActorSystemContext systemContext, TenantId tenantId) { |
|
||||
this.systemContext = systemContext; |
|
||||
this.ruleService = systemContext.getRuleService(); |
|
||||
this.ruleActors = new HashMap<>(); |
|
||||
this.tenantId = tenantId; |
|
||||
} |
|
||||
|
|
||||
public void init(ActorContext context) { |
|
||||
doInit(context); |
|
||||
} |
|
||||
|
|
||||
private void doInit(ActorContext context) { |
|
||||
PageDataIterable<RuleMetaData> ruleIterator = new PageDataIterable<>(getFetchRulesFunction(), |
|
||||
ContextAwareActor.ENTITY_PACK_LIMIT); |
|
||||
ruleMap = new HashMap<>(); |
|
||||
|
|
||||
for (RuleMetaData rule : ruleIterator) { |
|
||||
log.debug("[{}] Creating rule actor {}", rule.getId(), rule); |
|
||||
ActorRef ref = getOrCreateRuleActor(context, rule.getId()); |
|
||||
ruleMap.put(rule, RuleActorMetaData.systemRule(rule.getId(), rule.getWeight(), ref)); |
|
||||
log.debug("[{}] Rule actor created.", rule.getId()); |
|
||||
} |
|
||||
|
|
||||
refreshRuleChain(); |
|
||||
} |
|
||||
|
|
||||
public Optional<ActorRef> update(ActorContext context, RuleId ruleId, ComponentLifecycleEvent event) { |
|
||||
if (ruleMap == null) { |
|
||||
doInit(context); |
|
||||
} |
|
||||
RuleMetaData rule; |
|
||||
if (event != ComponentLifecycleEvent.DELETED) { |
|
||||
rule = systemContext.getRuleService().findRuleById(ruleId); |
|
||||
} else { |
|
||||
rule = ruleMap.keySet().stream() |
|
||||
.filter(r -> r.getId().equals(ruleId)) |
|
||||
.peek(r -> r.setState(ComponentLifecycleState.SUSPENDED)) |
|
||||
.findFirst() |
|
||||
.orElse(null); |
|
||||
if (rule != null) { |
|
||||
ruleMap.remove(rule); |
|
||||
ruleActors.remove(ruleId); |
|
||||
} |
|
||||
} |
|
||||
if (rule != null) { |
|
||||
RuleActorMetaData actorMd = ruleMap.get(rule); |
|
||||
if (actorMd == null) { |
|
||||
ActorRef ref = getOrCreateRuleActor(context, rule.getId()); |
|
||||
actorMd = RuleActorMetaData.systemRule(rule.getId(), rule.getWeight(), ref); |
|
||||
ruleMap.put(rule, actorMd); |
|
||||
} |
|
||||
refreshRuleChain(); |
|
||||
return Optional.of(actorMd.getActorRef()); |
|
||||
} else { |
|
||||
log.warn("[{}] Can't process unknown rule!", ruleId); |
|
||||
return Optional.empty(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
abstract FetchFunction<RuleMetaData> getFetchRulesFunction(); |
|
||||
|
|
||||
abstract String getDispatcherName(); |
|
||||
|
|
||||
public ActorRef getOrCreateRuleActor(ActorContext context, RuleId ruleId) { |
|
||||
return ruleActors.computeIfAbsent(ruleId, rId -> |
|
||||
context.actorOf(Props.create(new RuleActor.ActorCreator(systemContext, tenantId, rId)) |
|
||||
.withDispatcher(getDispatcherName()), rId.toString())); |
|
||||
} |
|
||||
|
|
||||
public RuleActorChain getRuleChain(ActorContext context) { |
|
||||
if (ruleChain == null) { |
|
||||
doInit(context); |
|
||||
} |
|
||||
return ruleChain; |
|
||||
} |
|
||||
|
|
||||
private void refreshRuleChain() { |
|
||||
Set<RuleActorMetaData> activeRuleSet = new HashSet<>(); |
|
||||
for (Map.Entry<RuleMetaData, RuleActorMetaData> rule : ruleMap.entrySet()) { |
|
||||
if (rule.getKey().getState() == ComponentLifecycleState.ACTIVE) { |
|
||||
activeRuleSet.add(rule.getValue()); |
|
||||
} |
|
||||
} |
|
||||
ruleChain = new SimpleRuleActorChain(activeRuleSet); |
|
||||
} |
|
||||
} |
|
||||
@ -1,40 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.shared.rule; |
|
||||
|
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.service.DefaultActorService; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.PageDataIterable.FetchFunction; |
|
||||
import org.thingsboard.server.common.data.rule.RuleMetaData; |
|
||||
import org.thingsboard.server.dao.model.ModelConstants; |
|
||||
|
|
||||
public class SystemRuleManager extends RuleManager { |
|
||||
|
|
||||
public SystemRuleManager(ActorSystemContext systemContext) { |
|
||||
super(systemContext, new TenantId(ModelConstants.NULL_UUID)); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
FetchFunction<RuleMetaData> getFetchRulesFunction() { |
|
||||
return ruleService::findSystemRules; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
String getDispatcherName() { |
|
||||
return DefaultActorService.SYSTEM_RULE_DISPATCHER_NAME; |
|
||||
} |
|
||||
} |
|
||||
@ -1,48 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.shared.rule; |
|
||||
|
|
||||
import akka.actor.ActorContext; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.service.DefaultActorService; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.PageDataIterable.FetchFunction; |
|
||||
import org.thingsboard.server.common.data.rule.RuleMetaData; |
|
||||
|
|
||||
public class TenantRuleManager extends RuleManager { |
|
||||
|
|
||||
public TenantRuleManager(ActorSystemContext systemContext, TenantId tenantId) { |
|
||||
super(systemContext, tenantId); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void init(ActorContext context) { |
|
||||
if (systemContext.isTenantComponentsInitEnabled()) { |
|
||||
super.init(context); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
FetchFunction<RuleMetaData> getFetchRulesFunction() { |
|
||||
return link -> ruleService.findTenantRules(tenantId, link); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
String getDispatcherName() { |
|
||||
return DefaultActorService.TENANT_RULE_DISPATCHER_NAME; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,40 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.tenant; |
|
||||
|
|
||||
import org.thingsboard.server.actors.rule.RuleActorChain; |
|
||||
import org.thingsboard.server.common.msg.device.ToDeviceActorMsg; |
|
||||
|
|
||||
public class RuleChainDeviceMsg { |
|
||||
|
|
||||
private final ToDeviceActorMsg toDeviceActorMsg; |
|
||||
private final RuleActorChain ruleChain; |
|
||||
|
|
||||
public RuleChainDeviceMsg(ToDeviceActorMsg toDeviceActorMsg, RuleActorChain ruleChain) { |
|
||||
super(); |
|
||||
this.toDeviceActorMsg = toDeviceActorMsg; |
|
||||
this.ruleChain = ruleChain; |
|
||||
} |
|
||||
|
|
||||
public ToDeviceActorMsg getToDeviceActorMsg() { |
|
||||
return toDeviceActorMsg; |
|
||||
} |
|
||||
|
|
||||
public RuleActorChain getRuleChain() { |
|
||||
return ruleChain; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -0,0 +1,41 @@ |
|||||
|
package org.thingsboard.server.controller; |
||||
|
|
||||
|
import com.fasterxml.jackson.core.type.TypeReference; |
||||
|
import org.thingsboard.server.common.data.DataConstants; |
||||
|
import org.thingsboard.server.common.data.Event; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.RuleChainId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.page.TimePageData; |
||||
|
import org.thingsboard.server.common.data.page.TimePageLink; |
||||
|
import org.thingsboard.server.common.data.rule.RuleChain; |
||||
|
import org.thingsboard.server.common.data.rule.RuleChainMetaData; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 20.03.18. |
||||
|
*/ |
||||
|
public class AbstractRuleEngineControllerTest extends AbstractControllerTest{ |
||||
|
|
||||
|
protected RuleChain saveRuleChain(RuleChain ruleChain) throws Exception { |
||||
|
return doPost("/api/ruleChain", ruleChain, RuleChain.class); |
||||
|
} |
||||
|
|
||||
|
protected RuleChain getRuleChain(RuleChainId ruleChainId) throws Exception { |
||||
|
return doGet("/api/ruleChain/" + ruleChainId.getId().toString(), RuleChain.class); |
||||
|
} |
||||
|
|
||||
|
protected RuleChainMetaData saveRuleChainMetaData(RuleChainMetaData ruleChainMD) throws Exception { |
||||
|
return doPost("/api/ruleChain/metadata", ruleChainMD, RuleChainMetaData.class); |
||||
|
} |
||||
|
|
||||
|
protected RuleChainMetaData getRuleChainMetaData(RuleChainId ruleChainId) throws Exception { |
||||
|
return doGet("/api/ruleChain/metadata/" + ruleChainId.getId().toString(), RuleChainMetaData.class); |
||||
|
} |
||||
|
|
||||
|
protected TimePageData<Event> getDebugEvents(TenantId tenantId, EntityId entityId, int limit) throws Exception { |
||||
|
TimePageLink pageLink = new TimePageLink(limit); |
||||
|
return doGetTypedWithTimePageLink("/api/events/{entityType}/{entityId}/{eventType}?tenantId={tenantId}&", |
||||
|
new TypeReference<TimePageData<Event>>() { |
||||
|
}, pageLink, entityId.getEntityType(), entityId.getId(), DataConstants.DEBUG, tenantId.getId()); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,35 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.rules; |
||||
|
|
||||
|
import org.junit.ClassRule; |
||||
|
import org.junit.extensions.cpsuite.ClasspathSuite; |
||||
|
import org.junit.runner.RunWith; |
||||
|
import org.thingsboard.server.dao.CustomSqlUnit; |
||||
|
|
||||
|
import java.util.Arrays; |
||||
|
|
||||
|
@RunWith(ClasspathSuite.class) |
||||
|
@ClasspathSuite.ClassnameFilters({ |
||||
|
"org.thingsboard.server.rules.flow.*Test", "org.thingsboard.server.rules.lifecycle.*Test"}) |
||||
|
public class RuleEngineSqlTestSuite { |
||||
|
|
||||
|
@ClassRule |
||||
|
public static CustomSqlUnit sqlUnit = new CustomSqlUnit( |
||||
|
Arrays.asList("sql/schema.sql", "sql/system-data.sql"), |
||||
|
"sql/drop-all-tables.sql", |
||||
|
"sql-test.properties"); |
||||
|
} |
||||
@ -0,0 +1,156 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 The Thingsboard Authors |
||||
|
* <p> |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* <p> |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* <p> |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.rules.flow; |
||||
|
|
||||
|
import com.datastax.driver.core.utils.UUIDs; |
||||
|
import lombok.Data; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.junit.After; |
||||
|
import org.junit.Assert; |
||||
|
import org.junit.Before; |
||||
|
import org.junit.Test; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.thingsboard.rule.engine.metadata.TbGetAttributesNodeConfiguration; |
||||
|
import org.thingsboard.server.actors.service.ActorService; |
||||
|
import org.thingsboard.server.common.data.*; |
||||
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
||||
|
import org.thingsboard.server.common.data.page.TimePageData; |
||||
|
import org.thingsboard.server.common.data.rule.RuleChain; |
||||
|
import org.thingsboard.server.common.data.rule.RuleChainMetaData; |
||||
|
import org.thingsboard.server.common.data.rule.RuleNode; |
||||
|
import org.thingsboard.server.common.data.security.Authority; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
||||
|
import org.thingsboard.server.common.msg.system.ServiceToRuleEngineMsg; |
||||
|
import org.thingsboard.server.controller.AbstractRuleEngineControllerTest; |
||||
|
import org.thingsboard.server.dao.attributes.AttributesService; |
||||
|
|
||||
|
import java.util.Collections; |
||||
|
|
||||
|
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; |
||||
|
|
||||
|
/** |
||||
|
* @author Valerii Sosliuk |
||||
|
*/ |
||||
|
@Slf4j |
||||
|
public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRuleEngineControllerTest { |
||||
|
|
||||
|
private static final String MQTT_URL = "tcp://localhost:1883"; |
||||
|
private static final Long TIME_TO_HANDLE_REQUEST = 500L; |
||||
|
|
||||
|
private Tenant savedTenant; |
||||
|
private User tenantAdmin; |
||||
|
|
||||
|
@Autowired |
||||
|
private ActorService actorService; |
||||
|
|
||||
|
@Autowired |
||||
|
private AttributesService attributesService; |
||||
|
|
||||
|
@Before |
||||
|
public void beforeTest() throws Exception { |
||||
|
loginSysAdmin(); |
||||
|
|
||||
|
Tenant tenant = new Tenant(); |
||||
|
tenant.setTitle("My tenant"); |
||||
|
savedTenant = doPost("/api/tenant", tenant, Tenant.class); |
||||
|
Assert.assertNotNull(savedTenant); |
||||
|
|
||||
|
tenantAdmin = new User(); |
||||
|
tenantAdmin.setAuthority(Authority.TENANT_ADMIN); |
||||
|
tenantAdmin.setTenantId(savedTenant.getId()); |
||||
|
tenantAdmin.setEmail("tenant2@thingsboard.org"); |
||||
|
tenantAdmin.setFirstName("Joe"); |
||||
|
tenantAdmin.setLastName("Downs"); |
||||
|
|
||||
|
createUserAndLogin(tenantAdmin, "testPassword1"); |
||||
|
} |
||||
|
|
||||
|
@After |
||||
|
public void afterTest() throws Exception { |
||||
|
loginSysAdmin(); |
||||
|
if (savedTenant != null) { |
||||
|
doDelete("/api/tenant/" + savedTenant.getId().getId().toString()).andExpect(status().isOk()); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testSimpleRuleChainCreation() throws Exception { |
||||
|
// Creating Rule Chain
|
||||
|
RuleChain ruleChain = new RuleChain(); |
||||
|
ruleChain.setName("Simple Rule Chain"); |
||||
|
ruleChain.setTenantId(savedTenant.getId()); |
||||
|
ruleChain.setRoot(true); |
||||
|
ruleChain.setDebugMode(true); |
||||
|
ruleChain = saveRuleChain(ruleChain); |
||||
|
Assert.assertNull(ruleChain.getFirstRuleNodeId()); |
||||
|
|
||||
|
RuleChainMetaData metaData = new RuleChainMetaData(); |
||||
|
metaData.setRuleChainId(ruleChain.getId()); |
||||
|
|
||||
|
RuleNode ruleNode = new RuleNode(); |
||||
|
ruleNode.setName("Simple Rule Node"); |
||||
|
ruleNode.setType(org.thingsboard.rule.engine.metadata.TbGetAttributesNode.class.getName()); |
||||
|
ruleNode.setDebugMode(true); |
||||
|
TbGetAttributesNodeConfiguration configuration = new TbGetAttributesNodeConfiguration(); |
||||
|
configuration.setServerAttributeNames(Collections.singletonList("serverAttributeKey")); |
||||
|
ruleNode.setConfiguration(mapper.valueToTree(configuration)); |
||||
|
|
||||
|
metaData.setNodes(Collections.singletonList(ruleNode)); |
||||
|
metaData.setFirstNodeIndex(0); |
||||
|
|
||||
|
metaData = saveRuleChainMetaData(metaData); |
||||
|
Assert.assertNotNull(metaData); |
||||
|
|
||||
|
ruleChain = getRuleChain(ruleChain.getId()); |
||||
|
Assert.assertNotNull(ruleChain.getFirstRuleNodeId()); |
||||
|
|
||||
|
// Saving the device
|
||||
|
Device device = new Device(); |
||||
|
device.setName("My device"); |
||||
|
device.setType("default"); |
||||
|
device = doPost("/api/device", device, Device.class); |
||||
|
|
||||
|
attributesService.save(device.getId(), DataConstants.SERVER_SCOPE, |
||||
|
Collections.singletonList(new BaseAttributeKvEntry(new StringDataEntry("serverAttributeKey", "serverAttributeValue"), System.currentTimeMillis()))); |
||||
|
|
||||
|
// Pushing Message to the system
|
||||
|
TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), |
||||
|
"CUSTOM", |
||||
|
device.getId(), |
||||
|
new TbMsgMetaData(), |
||||
|
new byte[]{}); |
||||
|
actorService.onMsg(new ServiceToRuleEngineMsg(savedTenant.getId(), tbMsg)); |
||||
|
|
||||
|
Thread.sleep(3000); |
||||
|
|
||||
|
TimePageData<Event> events = getDebugEvents(savedTenant.getId(), ruleChain.getFirstRuleNodeId(), 1000); |
||||
|
|
||||
|
Assert.assertEquals(2, events.getData().size()); |
||||
|
|
||||
|
Event inEvent = events.getData().stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); |
||||
|
Assert.assertEquals(ruleChain.getFirstRuleNodeId(), inEvent.getEntityId()); |
||||
|
Assert.assertEquals(device.getId().getId().toString(), inEvent.getBody().get("entityId").asText()); |
||||
|
|
||||
|
Event outEvent = events.getData().stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.OUT)).findFirst().get(); |
||||
|
Assert.assertEquals(ruleChain.getFirstRuleNodeId(), outEvent.getEntityId()); |
||||
|
Assert.assertEquals(device.getId().getId().toString(), outEvent.getBody().get("entityId").asText()); |
||||
|
|
||||
|
Assert.assertEquals("serverAttributeValue", outEvent.getBody().get("metadata").get("ss.serverAttributeKey").asText()); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,43 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.api; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.plugin.ComponentScope; |
||||
|
import org.thingsboard.server.extensions.api.component.EmptyComponentConfiguration; |
||||
|
|
||||
|
import java.lang.annotation.ElementType; |
||||
|
import java.lang.annotation.Retention; |
||||
|
import java.lang.annotation.RetentionPolicy; |
||||
|
import java.lang.annotation.Target; |
||||
|
|
||||
|
/** |
||||
|
* @author Andrew Shvayka |
||||
|
*/ |
||||
|
@Retention(RetentionPolicy.RUNTIME) |
||||
|
@Target(ElementType.TYPE) |
||||
|
public @interface ActionNode { |
||||
|
|
||||
|
String name(); |
||||
|
|
||||
|
ComponentScope scope() default ComponentScope.TENANT; |
||||
|
|
||||
|
String descriptor() default "EmptyNodeDescriptor.json"; |
||||
|
|
||||
|
String[] relationTypes() default {"Success","Failure"}; |
||||
|
|
||||
|
boolean customRelations() default false; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,42 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.api; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.plugin.ComponentScope; |
||||
|
import org.thingsboard.server.extensions.api.component.EmptyComponentConfiguration; |
||||
|
|
||||
|
import java.lang.annotation.ElementType; |
||||
|
import java.lang.annotation.Retention; |
||||
|
import java.lang.annotation.RetentionPolicy; |
||||
|
import java.lang.annotation.Target; |
||||
|
|
||||
|
/** |
||||
|
* @author Andrew Shvayka |
||||
|
*/ |
||||
|
@Retention(RetentionPolicy.RUNTIME) |
||||
|
@Target(ElementType.TYPE) |
||||
|
public @interface EnrichmentNode { |
||||
|
|
||||
|
String name(); |
||||
|
|
||||
|
ComponentScope scope() default ComponentScope.TENANT; |
||||
|
|
||||
|
String descriptor() default "EmptyNodeDescriptor.json"; |
||||
|
|
||||
|
String[] relationTypes() default {"Success","Failure"}; |
||||
|
|
||||
|
boolean customRelations() default false; |
||||
|
} |
||||
@ -0,0 +1,43 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.api; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.plugin.ComponentScope; |
||||
|
import org.thingsboard.server.extensions.api.component.EmptyComponentConfiguration; |
||||
|
|
||||
|
import java.lang.annotation.ElementType; |
||||
|
import java.lang.annotation.Retention; |
||||
|
import java.lang.annotation.RetentionPolicy; |
||||
|
import java.lang.annotation.Target; |
||||
|
|
||||
|
/** |
||||
|
* @author Andrew Shvayka |
||||
|
*/ |
||||
|
@Retention(RetentionPolicy.RUNTIME) |
||||
|
@Target(ElementType.TYPE) |
||||
|
public @interface FilterNode { |
||||
|
|
||||
|
String name(); |
||||
|
|
||||
|
ComponentScope scope() default ComponentScope.TENANT; |
||||
|
|
||||
|
String descriptor() default "EmptyNodeDescriptor.json"; |
||||
|
|
||||
|
String[] relationTypes() default {"Success","Failure"}; |
||||
|
|
||||
|
boolean customRelations() default false; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,43 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2018 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.rule.engine.api; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.plugin.ComponentScope; |
||||
|
import org.thingsboard.server.extensions.api.component.EmptyComponentConfiguration; |
||||
|
|
||||
|
import java.lang.annotation.ElementType; |
||||
|
import java.lang.annotation.Retention; |
||||
|
import java.lang.annotation.RetentionPolicy; |
||||
|
import java.lang.annotation.Target; |
||||
|
|
||||
|
/** |
||||
|
* @author Andrew Shvayka |
||||
|
*/ |
||||
|
@Retention(RetentionPolicy.RUNTIME) |
||||
|
@Target(ElementType.TYPE) |
||||
|
public @interface TransformationNode { |
||||
|
|
||||
|
String name(); |
||||
|
|
||||
|
ComponentScope scope() default ComponentScope.TENANT; |
||||
|
|
||||
|
String descriptor() default "EmptyNodeDescriptor.json"; |
||||
|
|
||||
|
String[] relationTypes() default {"Success","Failure"}; |
||||
|
|
||||
|
boolean customRelations() default false; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,2 @@ |
|||||
|
{ |
||||
|
} |
||||
Loading…
Reference in new issue