Browse Source

Fix conflicts

pull/4772/head
Igor Kulikov 5 years ago
parent
commit
48af7d37c3
  1. 14
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  2. 6
      application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java
  3. 158
      application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java
  4. 15
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java
  5. 4
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  6. 30
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java
  7. 50
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java
  8. 29
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java
  9. 19
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java

14
application/src/main/java/org/thingsboard/server/controller/RuleChainController.java

@ -74,6 +74,7 @@ import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Set;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@Slf4j @Slf4j
@ -86,6 +87,7 @@ public class RuleChainController extends BaseController {
public static final String RULE_NODE_ID = "ruleNodeId"; public static final String RULE_NODE_ID = "ruleNodeId";
private static final ObjectMapper objectMapper = new ObjectMapper(); private static final ObjectMapper objectMapper = new ObjectMapper();
public static final int TIMEOUT = 20;
@Autowired @Autowired
private InstallScripts installScripts; private InstallScripts installScripts;
@ -388,25 +390,25 @@ public class RuleChainController extends BaseController {
TbMsg inMsg = TbMsg.newMsg(msgType, null, new TbMsgMetaData(metadata), TbMsgDataType.JSON, data); TbMsg inMsg = TbMsg.newMsg(msgType, null, new TbMsgMetaData(metadata), TbMsgDataType.JSON, data);
switch (scriptType) { switch (scriptType) {
case "update": case "update":
output = msgToOutput(engine.executeUpdate(inMsg)); output = msgToOutput(engine.executeUpdateAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS));
break; break;
case "generate": case "generate":
output = msgToOutput(engine.executeGenerate(inMsg)); output = msgToOutput(engine.executeGenerateAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS));
break; break;
case "filter": case "filter":
boolean result = engine.executeFilter(inMsg); boolean result = engine.executeFilterAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS);
output = Boolean.toString(result); output = Boolean.toString(result);
break; break;
case "switch": case "switch":
Set<String> states = engine.executeSwitch(inMsg); Set<String> states = engine.executeSwitchAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS);
output = objectMapper.writeValueAsString(states); output = objectMapper.writeValueAsString(states);
break; break;
case "json": case "json":
JsonNode json = engine.executeJson(inMsg); JsonNode json = engine.executeJsonAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS);
output = objectMapper.writeValueAsString(json); output = objectMapper.writeValueAsString(json);
break; break;
case "string": case "string":
output = engine.executeToString(inMsg); output = engine.executeToStringAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS);
break; break;
default: default:
throw new IllegalArgumentException("Unsupported script type: " + scriptType); throw new IllegalArgumentException("Unsupported script type: " + scriptType);

6
application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java

@ -25,6 +25,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.scheduling.annotation.Scheduled; import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.util.StopWatch;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.queue.TbQueueRequestTemplate;
@ -180,6 +181,9 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
.setInvokeRequest(jsRequestBuilder.build()) .setInvokeRequest(jsRequestBuilder.build())
.build(); .build();
StopWatch stopWatch = new StopWatch();
stopWatch.start();
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper));
if (maxRequestsTimeout > 0) { if (maxRequestsTimeout > 0) {
future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService);
@ -201,6 +205,8 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
} }
}, callbackExecutor); }, callbackExecutor);
return Futures.transform(future, response -> { return Futures.transform(future, response -> {
stopWatch.stop();
log.trace("doInvokeFunction js-response took {}ms for uuid {}", stopWatch.getTotalTimeMillis(), response.getKey());
JsInvokeProtos.JsInvokeResponse invokeResult = response.getValue().getInvokeResponse(); JsInvokeProtos.JsInvokeResponse invokeResult = response.getValue().getInvokeResponse();
if (invokeResult.getSuccess()) { if (invokeResult.getSuccess()) {
return invokeResult.getResult(); return invokeResult.getResult();

158
application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java

@ -18,12 +18,12 @@ package org.thingsboard.server.service.script;
import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.collect.Sets;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
@ -32,6 +32,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import javax.script.ScriptException; import javax.script.ScriptException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collections; import java.util.Collections;
import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Set;
@ -102,140 +103,115 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S
String newMessageType = !StringUtils.isEmpty(messageType) ? messageType : msg.getType(); String newMessageType = !StringUtils.isEmpty(messageType) ? messageType : msg.getType();
return TbMsg.transformMsg(msg, newMessageType, msg.getOriginator(), newMetadata, newData); return TbMsg.transformMsg(msg, newMessageType, msg.getOriginator(), newMetadata, newData);
} catch (Throwable th) { } catch (Throwable th) {
th.printStackTrace();
throw new RuntimeException("Failed to unbind message data from javascript result", th); throw new RuntimeException("Failed to unbind message data from javascript result", th);
} }
} }
@Override @Override
public List<TbMsg> executeUpdate(TbMsg msg) throws ScriptException { public ListenableFuture<List<TbMsg>> executeUpdateAsync(TbMsg msg) {
JsonNode result = executeScript(msg); ListenableFuture<JsonNode> result = executeScriptAsync(msg);
if (result.isObject()) { return Futures.transformAsync(result,
return Collections.singletonList(unbindMsg(result, msg)); json -> executeUpdateTransform(msg, json),
} else if (result.isArray()){ MoreExecutors.directExecutor());
List<TbMsg> res = new ArrayList<>(result.size()); }
result.forEach(jsonObject -> res.add(unbindMsg(jsonObject, msg)));
return res; ListenableFuture<List<TbMsg>> executeUpdateTransform(TbMsg msg, JsonNode json) {
} else { if (json.isObject()) {
log.warn("Wrong result type: {}", result.getNodeType()); return Futures.immediateFuture(Collections.singletonList(unbindMsg(json, msg)));
throw new ScriptException("Wrong result type: " + result.getNodeType()); } else if (json.isArray()) {
List<TbMsg> res = new ArrayList<>(json.size());
json.forEach(jsonObject -> res.add(unbindMsg(jsonObject, msg)));
return Futures.immediateFuture(res);
} }
log.warn("Wrong result type: {}", json.getNodeType());
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + json.getNodeType()));
} }
@Override @Override
public ListenableFuture<List<TbMsg>> executeUpdateAsync(TbMsg msg) { public ListenableFuture<TbMsg> executeGenerateAsync(TbMsg prevMsg) {
ListenableFuture<JsonNode> result = executeScriptAsync(msg); return Futures.transformAsync(executeScriptAsync(prevMsg),
return Futures.transformAsync(result, json -> { result -> executeGenerateTransform(prevMsg, result),
if (json.isObject()) { MoreExecutors.directExecutor());
return Futures.immediateFuture(Collections.singletonList(unbindMsg(json, msg)));
} else if (json.isArray()){
List<TbMsg> res = new ArrayList<>(json.size());
json.forEach(jsonObject -> res.add(unbindMsg(jsonObject, msg)));
return Futures.immediateFuture(res);
}
else{
log.warn("Wrong result type: {}", json.getNodeType());
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + json.getNodeType()));
}
}, MoreExecutors.directExecutor());
} }
@Override ListenableFuture<TbMsg> executeGenerateTransform(TbMsg prevMsg, JsonNode result) {
public TbMsg executeGenerate(TbMsg prevMsg) throws ScriptException {
JsonNode result = executeScript(prevMsg);
if (!result.isObject()) { if (!result.isObject()) {
log.warn("Wrong result type: {}", result.getNodeType()); log.warn("Wrong result type: {}", result.getNodeType());
throw new ScriptException("Wrong result type: " + result.getNodeType()); Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + result.getNodeType()));
} }
return unbindMsg(result, prevMsg); return Futures.immediateFuture(unbindMsg(result, prevMsg));
} }
@Override @Override
public JsonNode executeJson(TbMsg msg) throws ScriptException { public ListenableFuture<JsonNode> executeJsonAsync(TbMsg msg) {
return executeScript(msg);
}
@Override
public ListenableFuture<JsonNode> executeJsonAsync(TbMsg msg) throws ScriptException {
return executeScriptAsync(msg); return executeScriptAsync(msg);
} }
@Override @Override
public String executeToString(TbMsg msg) throws ScriptException { public ListenableFuture<String> executeToStringAsync(TbMsg msg) {
JsonNode result = executeScript(msg); return Futures.transformAsync(executeScriptAsync(msg),
if (!result.isTextual()) { this::executeToStringTransform,
log.warn("Wrong result type: {}", result.getNodeType()); MoreExecutors.directExecutor());
throw new ScriptException("Wrong result type: " + result.getNodeType());
}
return result.asText();
} }
@Override ListenableFuture<String> executeToStringTransform(JsonNode result) {
public boolean executeFilter(TbMsg msg) throws ScriptException { if (result.isTextual()) {
JsonNode result = executeScript(msg); return Futures.immediateFuture(result.asText());
if (!result.isBoolean()) {
log.warn("Wrong result type: {}", result.getNodeType());
throw new ScriptException("Wrong result type: " + result.getNodeType());
} }
return result.asBoolean(); log.warn("Wrong result type: {}", result.getNodeType());
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + result.getNodeType()));
} }
@Override @Override
public ListenableFuture<Boolean> executeFilterAsync(TbMsg msg) { public ListenableFuture<Boolean> executeFilterAsync(TbMsg msg) {
ListenableFuture<JsonNode> result = executeScriptAsync(msg); return Futures.transformAsync(executeScriptAsync(msg),
return Futures.transformAsync(result, json -> { this::executeFilterTransform,
if (!json.isBoolean()) { MoreExecutors.directExecutor());
log.warn("Wrong result type: {}", json.getNodeType());
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + json.getNodeType()));
} else {
return Futures.immediateFuture(json.asBoolean());
}
}, MoreExecutors.directExecutor());
} }
@Override ListenableFuture<Boolean> executeFilterTransform(JsonNode json) {
public Set<String> executeSwitch(TbMsg msg) throws ScriptException { if (json.isBoolean()) {
JsonNode result = executeScript(msg); return Futures.immediateFuture(json.asBoolean());
}
log.warn("Wrong result type: {}", json.getNodeType());
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + json.getNodeType()));
}
ListenableFuture<Set<String>> executeSwitchTransform(JsonNode result) {
if (result.isTextual()) { if (result.isTextual()) {
return Collections.singleton(result.asText()); return Futures.immediateFuture(Collections.singleton(result.asText()));
} else if (result.isArray()) { }
Set<String> nextStates = Sets.newHashSet(); if (result.isArray()) {
Set<String> nextStates = new HashSet<>();
for (JsonNode val : result) { for (JsonNode val : result) {
if (!val.isTextual()) { if (!val.isTextual()) {
log.warn("Wrong result type: {}", val.getNodeType()); log.warn("Wrong result type: {}", val.getNodeType());
throw new ScriptException("Wrong result type: " + val.getNodeType()); return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + val.getNodeType()));
} else { } else {
nextStates.add(val.asText()); nextStates.add(val.asText());
} }
} }
return nextStates; return Futures.immediateFuture(nextStates);
} else {
log.warn("Wrong result type: {}", result.getNodeType());
throw new ScriptException("Wrong result type: " + result.getNodeType());
} }
log.warn("Wrong result type: {}", result.getNodeType());
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + result.getNodeType()));
} }
private JsonNode executeScript(TbMsg msg) throws ScriptException { @Override
try { public ListenableFuture<Set<String>> executeSwitchAsync(TbMsg msg) {
String[] inArgs = prepareArgs(msg); return Futures.transformAsync(executeScriptAsync(msg),
String eval = sandboxService.invokeFunction(tenantId, msg.getCustomerId(), this.scriptId, inArgs[0], inArgs[1], inArgs[2]).get().toString(); this::executeSwitchTransform,
return mapper.readTree(eval); MoreExecutors.directExecutor()); //usually runs in a callbackExecutor
} catch (ExecutionException e) {
if (e.getCause() instanceof ScriptException) {
throw (ScriptException) e.getCause();
} else if (e.getCause() instanceof RuntimeException) {
throw new ScriptException(e.getCause().getMessage());
} else {
throw new ScriptException(e);
}
} catch (Exception e) {
throw new ScriptException(e);
}
} }
private ListenableFuture<JsonNode> executeScriptAsync(TbMsg msg) { ListenableFuture<JsonNode> executeScriptAsync(TbMsg msg) {
log.trace("execute script async, msg {}", msg);
String[] inArgs = prepareArgs(msg); String[] inArgs = prepareArgs(msg);
return Futures.transformAsync(sandboxService.invokeFunction(tenantId, msg.getCustomerId(), this.scriptId, inArgs[0], inArgs[1], inArgs[2]), return executeScriptAsync(msg.getCustomerId(), inArgs[0], inArgs[1], inArgs[2]);
}
ListenableFuture<JsonNode> executeScriptAsync(CustomerId customerId, Object... args) {
return Futures.transformAsync(sandboxService.invokeFunction(tenantId, customerId, this.scriptId, args),
o -> { o -> {
try { try {
return Futures.immediateFuture(mapper.readTree(o.toString())); return Futures.immediateFuture(mapper.readTree(o.toString()));

15
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java

@ -19,29 +19,22 @@ import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import javax.script.ScriptException;
import java.util.List; import java.util.List;
import java.util.Set; import java.util.Set;
public interface ScriptEngine { public interface ScriptEngine {
List<TbMsg> executeUpdate(TbMsg msg) throws ScriptException;
ListenableFuture<List<TbMsg>> executeUpdateAsync(TbMsg msg); ListenableFuture<List<TbMsg>> executeUpdateAsync(TbMsg msg);
TbMsg executeGenerate(TbMsg prevMsg) throws ScriptException; ListenableFuture<TbMsg> executeGenerateAsync(TbMsg prevMsg);
boolean executeFilter(TbMsg msg) throws ScriptException;
ListenableFuture<Boolean> executeFilterAsync(TbMsg msg); ListenableFuture<Boolean> executeFilterAsync(TbMsg msg);
Set<String> executeSwitch(TbMsg msg) throws ScriptException; ListenableFuture<Set<String>> executeSwitchAsync(TbMsg msg);
JsonNode executeJson(TbMsg msg) throws ScriptException;
ListenableFuture<JsonNode> executeJsonAsync(TbMsg msg) throws ScriptException; ListenableFuture<JsonNode> executeJsonAsync(TbMsg msg);
String executeToString(TbMsg msg) throws ScriptException; ListenableFuture<String> executeToStringAsync(TbMsg msg);
void destroy(); void destroy();

4
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -214,6 +214,10 @@ public interface TbContext {
EdgeEventService getEdgeEventService(); EdgeEventService getEdgeEventService();
/**
* Js script executors call are completely asynchronous
* */
@Deprecated
ListeningExecutor getJsExecutor(); ListeningExecutor getJsExecutor();
ListeningExecutor getMailExecutor(); ListeningExecutor getMailExecutor();

30
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java

@ -15,8 +15,11 @@
*/ */
package org.thingsboard.rule.engine.action; package org.thingsboard.rule.engine.action;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.ListeningExecutor; import org.checkerframework.checker.nullness.qual.Nullable;
import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.ScriptEngine; import org.thingsboard.rule.engine.api.ScriptEngine;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
@ -55,18 +58,21 @@ public class TbLogNode implements TbNode {
@Override @Override
public void onMsg(TbContext ctx, TbMsg msg) { public void onMsg(TbContext ctx, TbMsg msg) {
ListeningExecutor jsExecutor = ctx.getJsExecutor();
ctx.logJsEvalRequest(); ctx.logJsEvalRequest();
withCallback(jsExecutor.executeAsync(() -> jsEngine.executeToString(msg)), Futures.addCallback(jsEngine.executeToStringAsync(msg), new FutureCallback<String>() {
toString -> { @Override
ctx.logJsEvalResponse(); public void onSuccess(@Nullable String result) {
log.info(toString); ctx.logJsEvalResponse();
ctx.tellSuccess(msg); log.info(result);
}, ctx.tellSuccess(msg);
t -> { }
ctx.logJsEvalResponse();
ctx.tellFailure(msg, t); @Override
}); public void onFailure(Throwable t) {
ctx.logJsEvalResponse();
ctx.tellFailure(msg, t);
}
}, MoreExecutors.directExecutor()); //usually js responses runs on js callback executor
} }
@Override @Override

50
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java

@ -15,9 +15,12 @@
*/ */
package org.thingsboard.rule.engine.debug; package org.thingsboard.rule.engine.debug;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;
import org.thingsboard.common.util.TbStopWatch;
import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.ScriptEngine; import org.thingsboard.rule.engine.api.ScriptEngine;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
@ -35,6 +38,7 @@ import org.thingsboard.server.common.msg.queue.ServiceQueue;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import static org.thingsboard.common.util.DonAsynchron.withCallback; import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
@ -64,10 +68,11 @@ public class TbMsgGeneratorNode implements TbNode {
private EntityId originatorId; private EntityId originatorId;
private UUID nextTickId; private UUID nextTickId;
private TbMsg prevMsg; private TbMsg prevMsg;
private volatile boolean initialized; private final AtomicBoolean initialized = new AtomicBoolean(false);
@Override @Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
log.trace("init generator with config {}", configuration);
this.config = TbNodeUtils.convert(configuration, TbMsgGeneratorNodeConfiguration.class); this.config = TbNodeUtils.convert(configuration, TbMsgGeneratorNodeConfiguration.class);
this.delay = TimeUnit.SECONDS.toMillis(config.getPeriodInSeconds()); this.delay = TimeUnit.SECONDS.toMillis(config.getPeriodInSeconds());
this.currentMsgCount = 0; this.currentMsgCount = 0;
@ -81,35 +86,39 @@ public class TbMsgGeneratorNode implements TbNode {
@Override @Override
public void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) { public void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) {
log.trace("onPartitionChangeMsg, PartitionChangeMsg {}, config {}", msg, config);
updateGeneratorState(ctx); updateGeneratorState(ctx);
} }
private void updateGeneratorState(TbContext ctx) { private void updateGeneratorState(TbContext ctx) {
log.trace("updateGeneratorState, config {}", config);
if (ctx.isLocalEntity(originatorId)) { if (ctx.isLocalEntity(originatorId)) {
if (!initialized) { if (initialized.compareAndSet(false, true)) {
initialized = true;
this.jsEngine = ctx.createJsScriptEngine(config.getJsScript(), "prevMsg", "prevMetadata", "prevMsgType"); this.jsEngine = ctx.createJsScriptEngine(config.getJsScript(), "prevMsg", "prevMetadata", "prevMsgType");
scheduleTickMsg(ctx); scheduleTickMsg(ctx);
} }
} else if (initialized) { } else if (initialized.compareAndSet(true, false)) {
initialized = false;
destroy(); destroy();
} }
} }
@Override @Override
public void onMsg(TbContext ctx, TbMsg msg) { public void onMsg(TbContext ctx, TbMsg msg) {
if (initialized && msg.getType().equals(TB_MSG_GENERATOR_NODE_MSG) && msg.getId().equals(nextTickId)) { log.trace("onMsg, config {}, msg {}", config, msg);
if (initialized.get() && msg.getType().equals(TB_MSG_GENERATOR_NODE_MSG) && msg.getId().equals(nextTickId)) {
TbStopWatch sw = TbStopWatch.startNew();
withCallback(generate(ctx, msg), withCallback(generate(ctx, msg),
m -> { m -> {
if (initialized && (config.getMsgCount() == TbMsgGeneratorNodeConfiguration.UNLIMITED_MSG_COUNT || currentMsgCount < config.getMsgCount())) { log.trace("onMsg onSuccess callback, took {}ms, config {}, msg {}", sw.stopAndGetTotalTimeMillis(), config, msg);
if (initialized.get() && (config.getMsgCount() == TbMsgGeneratorNodeConfiguration.UNLIMITED_MSG_COUNT || currentMsgCount < config.getMsgCount())) {
ctx.enqueueForTellNext(m, SUCCESS); ctx.enqueueForTellNext(m, SUCCESS);
scheduleTickMsg(ctx); scheduleTickMsg(ctx);
currentMsgCount++; currentMsgCount++;
} }
}, },
t -> { t -> {
if (initialized && (config.getMsgCount() == TbMsgGeneratorNodeConfiguration.UNLIMITED_MSG_COUNT || currentMsgCount < config.getMsgCount())) { log.warn("onMsg onFailure callback, took {}ms, config {}, msg {}, exception {}", sw.stopAndGetTotalTimeMillis(), config, msg, t);
if (initialized.get() && (config.getMsgCount() == TbMsgGeneratorNodeConfiguration.UNLIMITED_MSG_COUNT || currentMsgCount < config.getMsgCount())) {
ctx.tellFailure(msg, t); ctx.tellFailure(msg, t);
scheduleTickMsg(ctx); scheduleTickMsg(ctx);
currentMsgCount++; currentMsgCount++;
@ -119,6 +128,7 @@ public class TbMsgGeneratorNode implements TbNode {
} }
private void scheduleTickMsg(TbContext ctx) { private void scheduleTickMsg(TbContext ctx) {
log.trace("scheduleTickMsg, config {}", config);
long curTs = System.currentTimeMillis(); long curTs = System.currentTimeMillis();
if (lastScheduledTs == 0L) { if (lastScheduledTs == 0L) {
lastScheduledTs = curTs; lastScheduledTs = curTs;
@ -131,22 +141,26 @@ public class TbMsgGeneratorNode implements TbNode {
} }
private ListenableFuture<TbMsg> generate(TbContext ctx, TbMsg msg) { private ListenableFuture<TbMsg> generate(TbContext ctx, TbMsg msg) {
return ctx.getJsExecutor().executeAsync(() -> { log.trace("generate, config {}", config);
if (prevMsg == null) { if (prevMsg == null) {
prevMsg = ctx.newMsg(ServiceQueue.MAIN, "", originatorId, msg.getCustomerId(), new TbMsgMetaData(), "{}"); prevMsg = ctx.newMsg(ServiceQueue.MAIN, "", originatorId, msg.getCustomerId(), new TbMsgMetaData(), "{}");
} }
if (initialized) { if (initialized.get()) {
ctx.logJsEvalRequest(); ctx.logJsEvalRequest();
TbMsg generated = jsEngine.executeGenerate(prevMsg); return Futures.transformAsync(jsEngine.executeGenerateAsync(prevMsg), generated -> {
log.trace("generate process response, generated {}, config {}", generated, config);
ctx.logJsEvalResponse(); ctx.logJsEvalResponse();
prevMsg = ctx.newMsg(ServiceQueue.MAIN, generated.getType(), originatorId, msg.getCustomerId(), generated.getMetaData(), generated.getData()); prevMsg = ctx.newMsg(ServiceQueue.MAIN, generated.getType(), originatorId, msg.getCustomerId(), generated.getMetaData(), generated.getData());
} return Futures.immediateFuture(prevMsg);
return prevMsg; }, MoreExecutors.directExecutor()); //usually it runs on js-executor-remote-callback thread pool
}); }
return Futures.immediateFuture(prevMsg);
} }
@Override @Override
public void destroy() { public void destroy() {
log.trace("destroy, config {}", config);
prevMsg = null; prevMsg = null;
if (jsEngine != null) { if (jsEngine != null) {
jsEngine.destroy(); jsEngine.destroy();

29
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java

@ -15,7 +15,11 @@
*/ */
package org.thingsboard.rule.engine.filter; package org.thingsboard.rule.engine.filter;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.ScriptEngine; import org.thingsboard.rule.engine.api.ScriptEngine;
@ -29,8 +33,6 @@ import org.thingsboard.server.common.msg.TbMsg;
import java.util.Set; import java.util.Set;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
@Slf4j @Slf4j
@RuleNode( @RuleNode(
type = ComponentType.FILTER, type = ComponentType.FILTER,
@ -58,17 +60,20 @@ public class TbJsSwitchNode implements TbNode {
@Override @Override
public void onMsg(TbContext ctx, TbMsg msg) { public void onMsg(TbContext ctx, TbMsg msg) {
ListeningExecutor jsExecutor = ctx.getJsExecutor();
ctx.logJsEvalRequest(); ctx.logJsEvalRequest();
withCallback(jsExecutor.executeAsync(() -> jsEngine.executeSwitch(msg)), Futures.addCallback(jsEngine.executeSwitchAsync(msg), new FutureCallback<Set<String>>() {
result -> { @Override
ctx.logJsEvalResponse(); public void onSuccess(@Nullable Set<String> result) {
processSwitch(ctx, msg, result); ctx.logJsEvalResponse();
}, processSwitch(ctx, msg, result);
t -> { }
ctx.logJsEvalFailure();
ctx.tellFailure(msg, t); @Override
}, ctx.getDbCallbackExecutor()); public void onFailure(Throwable t) {
ctx.logJsEvalFailure();
ctx.tellFailure(msg, t);
}
}, MoreExecutors.directExecutor()); //usually runs in a callbackExecutor
} }
private void processSwitch(TbContext ctx, TbMsg msg, Set<String> nextRelations) { private void processSwitch(TbContext ctx, TbMsg msg, Set<String> nextRelations) {

19
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java

@ -64,7 +64,7 @@ public class TbJsSwitchNodeTest {
private RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); private RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased());
@Test @Test
public void multipleRoutesAreAllowed() throws TbNodeException, ScriptException { public void multipleRoutesAreAllowed() throws TbNodeException {
initWithScript(); initWithScript();
TbMsgMetaData metaData = new TbMsgMetaData(); TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("temp", "10"); metaData.putValue("temp", "10");
@ -72,11 +72,9 @@ public class TbJsSwitchNodeTest {
String rawJson = "{\"name\": \"Vit\", \"passed\": 5}"; String rawJson = "{\"name\": \"Vit\", \"passed\": 5}";
TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson, ruleChainId, ruleNodeId); TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson, ruleChainId, ruleNodeId);
mockJsExecutor(); when(scriptEngine.executeSwitchAsync(msg)).thenReturn(Futures.immediateFuture(Sets.newHashSet("one", "three")));
when(scriptEngine.executeSwitch(msg)).thenReturn(Sets.newHashSet("one", "three"));
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
verify(ctx).getJsExecutor();
verify(ctx).tellNext(msg, Sets.newHashSet("one", "three")); verify(ctx).tellNext(msg, Sets.newHashSet("one", "three"));
} }
@ -92,19 +90,6 @@ public class TbJsSwitchNodeTest {
node.init(ctx, nodeConfiguration); node.init(ctx, nodeConfiguration);
} }
@SuppressWarnings("unchecked")
private void mockJsExecutor() {
when(ctx.getJsExecutor()).thenReturn(executor);
doAnswer((Answer<ListenableFuture<Set<String>>>) invocationOnMock -> {
try {
Callable task = (Callable) (invocationOnMock.getArguments())[0];
return Futures.immediateFuture((Set<String>) task.call());
} catch (Throwable th) {
return Futures.immediateFailedFuture(th);
}
}).when(executor).executeAsync(ArgumentMatchers.any(Callable.class));
}
private void verifyError(TbMsg msg, String message, Class expectedClass) { private void verifyError(TbMsg msg, String message, Class expectedClass) {
ArgumentCaptor<Throwable> captor = ArgumentCaptor.forClass(Throwable.class); ArgumentCaptor<Throwable> captor = ArgumentCaptor.forClass(Throwable.class);
verify(ctx).tellFailure(same(msg), captor.capture()); verify(ctx).tellFailure(same(msg), captor.capture());

Loading…
Cancel
Save