From a3b31337cae7f0b639285855fd1cc050d05cfda4 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 5 May 2021 09:09:12 +0300 Subject: [PATCH 1/9] sync method replaced with async executeGenerateAsync for ScriptEngine api. affected Generator node, ruleChainController (merged with ce) --- .../controller/RuleChainController.java | 4 +- .../service/script/RemoteJsInvokeService.java | 8 ++- .../script/RuleNodeJsScriptEngine.java | 19 ++++--- .../rule/engine/api/ScriptEngine.java | 2 +- .../rule/engine/debug/TbMsgGeneratorNode.java | 50 ++++++++++++------- 5 files changed, 54 insertions(+), 29 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index d5fcc2706f..3b9af2f0f2 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/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.Set; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @Slf4j @@ -86,6 +87,7 @@ public class RuleChainController extends BaseController { public static final String RULE_NODE_ID = "ruleNodeId"; private static final ObjectMapper objectMapper = new ObjectMapper(); + public static final int TIMEOUT = 20; @Autowired private InstallScripts installScripts; @@ -391,7 +393,7 @@ public class RuleChainController extends BaseController { output = msgToOutput(engine.executeUpdate(inMsg)); break; case "generate": - output = msgToOutput(engine.executeGenerate(inMsg)); + output = msgToOutput(engine.executeGenerateAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS)); break; case "filter": boolean result = engine.executeFilter(inMsg); diff --git a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java index 334a471973..1bc9203333 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java @@ -18,7 +18,6 @@ package org.thingsboard.server.service.script; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -26,6 +25,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.springframework.util.StopWatch; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.queue.TbQueueRequestTemplate; @@ -161,6 +161,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { @Override protected ListenableFuture doInvokeFunction(UUID scriptId, String functionName, Object[] args) { + log.trace("doInvokeFunction js-request for uuid {} with timeout {}ms", scriptId, maxRequestsTimeout); String scriptBody = scriptIdToBodysMap.get(scriptId); if (scriptBody == null) { return Futures.immediateFailedFuture(new RuntimeException("No script body found for scriptId: [" + scriptId + "]!")); @@ -180,6 +181,9 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { .setInvokeRequest(jsRequestBuilder.build()) .build(); + StopWatch stopWatch = new StopWatch(); + stopWatch.start(); + ListenableFuture> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); if (maxRequestsTimeout > 0) { future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); @@ -201,6 +205,8 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { } }, callbackExecutor); 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(); if (invokeResult.getSuccess()) { return invokeResult.getResult(); diff --git a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java index 066ce71a58..c688d6978d 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java @@ -102,7 +102,6 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S String newMessageType = !StringUtils.isEmpty(messageType) ? messageType : msg.getType(); return TbMsg.transformMsg(msg, newMessageType, msg.getOriginator(), newMetadata, newData); } catch (Throwable th) { - th.printStackTrace(); throw new RuntimeException("Failed to unbind message data from javascript result", th); } } @@ -141,13 +140,16 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S } @Override - public TbMsg executeGenerate(TbMsg prevMsg) throws ScriptException { - JsonNode result = executeScript(prevMsg); - if (!result.isObject()) { - log.warn("Wrong result type: {}", result.getNodeType()); - throw new ScriptException("Wrong result type: " + result.getNodeType()); - } - return unbindMsg(result, prevMsg); + public ListenableFuture executeGenerateAsync(TbMsg prevMsg) { + log.trace("execute generate async, prevMsg {}", prevMsg); + return Futures.transformAsync(executeScriptAsync(prevMsg), result -> { + if (!result.isObject()) { + log.warn("Wrong result type: {}", result.getNodeType()); + throw new ScriptException("Wrong result type: " + result.getNodeType()); + } + return Futures.immediateFuture(unbindMsg(result, prevMsg)); + }, MoreExecutors.directExecutor()); + } @Override @@ -234,6 +236,7 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S } private ListenableFuture executeScriptAsync(TbMsg msg) { + log.trace("execute script async, msg {}", msg); String[] inArgs = prepareArgs(msg); return Futures.transformAsync(sandboxService.invokeFunction(tenantId, msg.getCustomerId(), this.scriptId, inArgs[0], inArgs[1], inArgs[2]), o -> { diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java index 49420e8feb..a5dcf535e4 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java @@ -29,7 +29,7 @@ public interface ScriptEngine { ListenableFuture> executeUpdateAsync(TbMsg msg); - TbMsg executeGenerate(TbMsg prevMsg) throws ScriptException; + ListenableFuture executeGenerateAsync(TbMsg prevMsg); boolean executeFilter(TbMsg msg) throws ScriptException; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java index b5a47b4a0d..8fa4b97063 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java +++ b/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; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; +import org.thingsboard.common.util.TbStopWatch; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.ScriptEngine; 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.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import static org.thingsboard.common.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; @@ -64,10 +68,11 @@ public class TbMsgGeneratorNode implements TbNode { private EntityId originatorId; private UUID nextTickId; private TbMsg prevMsg; - private volatile boolean initialized; + private final AtomicBoolean initialized = new AtomicBoolean(false); @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { + log.trace("init generator with config {}", configuration); this.config = TbNodeUtils.convert(configuration, TbMsgGeneratorNodeConfiguration.class); this.delay = TimeUnit.SECONDS.toMillis(config.getPeriodInSeconds()); this.currentMsgCount = 0; @@ -81,35 +86,39 @@ public class TbMsgGeneratorNode implements TbNode { @Override public void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) { + log.trace("onPartitionChangeMsg, PartitionChangeMsg {}, config {}", msg, config); updateGeneratorState(ctx); } private void updateGeneratorState(TbContext ctx) { + log.trace("updateGeneratorState, config {}", config); if (ctx.isLocalEntity(originatorId)) { - if (!initialized) { - initialized = true; + if (initialized.compareAndSet(false, true)) { this.jsEngine = ctx.createJsScriptEngine(config.getJsScript(), "prevMsg", "prevMetadata", "prevMsgType"); scheduleTickMsg(ctx); } - } else if (initialized) { - initialized = false; + } else if (initialized.compareAndSet(true, false)) { destroy(); } } @Override 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), 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); scheduleTickMsg(ctx); currentMsgCount++; } }, t -> { - if (initialized && (config.getMsgCount() == TbMsgGeneratorNodeConfiguration.UNLIMITED_MSG_COUNT || currentMsgCount < config.getMsgCount())) { + log.warn("onMsg onFailure callback, took {}ms, config {}, msg {}", sw.stopAndGetTotalTimeMillis(), config, msg); + if (initialized.get() && (config.getMsgCount() == TbMsgGeneratorNodeConfiguration.UNLIMITED_MSG_COUNT || currentMsgCount < config.getMsgCount())) { ctx.tellFailure(msg, t); scheduleTickMsg(ctx); currentMsgCount++; @@ -119,6 +128,7 @@ public class TbMsgGeneratorNode implements TbNode { } private void scheduleTickMsg(TbContext ctx) { + log.trace("scheduleTickMsg, config {}", config); long curTs = System.currentTimeMillis(); if (lastScheduledTs == 0L) { lastScheduledTs = curTs; @@ -131,22 +141,26 @@ public class TbMsgGeneratorNode implements TbNode { } private ListenableFuture generate(TbContext ctx, TbMsg msg) { - return ctx.getJsExecutor().executeAsync(() -> { - if (prevMsg == null) { - prevMsg = ctx.newMsg(ServiceQueue.MAIN, "", originatorId, msg.getCustomerId(), new TbMsgMetaData(), "{}"); - } - if (initialized) { - ctx.logJsEvalRequest(); - TbMsg generated = jsEngine.executeGenerate(prevMsg); + log.trace("generate, config {}", config); + if (prevMsg == null) { + prevMsg = ctx.newMsg(ServiceQueue.MAIN, "", originatorId, msg.getCustomerId(), new TbMsgMetaData(), "{}"); + } + if (initialized.get()) { + ctx.logJsEvalRequest(); + return Futures.transformAsync(jsEngine.executeGenerateAsync(prevMsg), generated -> { + log.trace("generate process response, generated {}, config {}", generated, config); ctx.logJsEvalResponse(); prevMsg = ctx.newMsg(ServiceQueue.MAIN, generated.getType(), originatorId, msg.getCustomerId(), generated.getMetaData(), generated.getData()); - } - return prevMsg; - }); + return Futures.immediateFuture(prevMsg); + }, MoreExecutors.directExecutor()); + } + return Futures.immediateFuture(prevMsg); + } @Override public void destroy() { + log.trace("destroy, config {}", config); prevMsg = null; if (jsEngine != null) { jsEngine.destroy(); From 6888a1b349de234db1f6f3a8b81326fb34f438b1 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 5 May 2021 08:22:38 +0300 Subject: [PATCH 2/9] new utility class tbStopWatch added to improve performance measurement experience --- common/util/pom.xml | 5 ++ .../thingsboard/common/util/TbStopWatch.java | 62 +++++++++++++++++++ 2 files changed, 67 insertions(+) create mode 100644 common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java diff --git a/common/util/pom.xml b/common/util/pom.xml index 6f66da5790..128bc478d5 100644 --- a/common/util/pom.xml +++ b/common/util/pom.xml @@ -36,6 +36,11 @@ + + org.springframework + spring-core + ${spring.version} + com.google.guava guava diff --git a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java new file mode 100644 index 0000000000..b7dbd0052c --- /dev/null +++ b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java @@ -0,0 +1,62 @@ +/** + * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * + * Copyright © 2016-2021 ThingsBoard, Inc. All Rights Reserved. + * + * NOTICE: All information contained herein is, and remains + * the property of ThingsBoard, Inc. and its suppliers, + * if any. The intellectual and technical concepts contained + * herein are proprietary to ThingsBoard, Inc. + * and its suppliers and may be covered by U.S. and Foreign Patents, + * patents in process, and are protected by trade secret or copyright law. + * + * Dissemination of this information or reproduction of this material is strictly forbidden + * unless prior written permission is obtained from COMPANY. + * + * Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees, + * managers or contractors who have executed Confidentiality and Non-disclosure agreements + * explicitly covering such access. + * + * The copyright notice above does not evidence any actual or intended publication + * or disclosure of this source code, which includes + * information that is confidential and/or proprietary, and is a trade secret, of COMPANY. + * ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE, + * OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT + * THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED, + * AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES. + * THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION + * DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS, + * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. + */ +package org.thingsboard.common.util; + +import org.springframework.util.StopWatch; + +/** + * Utility method that extends Spring Framework StopWatch + * It is a MONOTONIC time stopwatch. + * It is a replacement for any measurements with a wall-clock like System.currentTimeMillis() + * It is not affected by leap second, day-light saving and wall-clock adjustments by manual or network time synchronization + * The main features is a single call for common use cases: + * - create and start: TbStopWatch sw = TbStopWatch.startNew() + * - stop and get: sw.stopAndGetTotalTimeMillis() or sw.stopAndGetLastTaskTimeMillis() + * */ +public class TbStopWatch extends StopWatch { + + public static TbStopWatch startNew(){ + TbStopWatch stopWatch = new TbStopWatch(); + stopWatch.start(); + return stopWatch; + } + + public long stopAndGetTotalTimeMillis(){ + stop(); + return getTotalTimeMillis(); + } + + public long stopAndGetLastTaskTimeMillis(){ + stop(); + return getLastTaskTimeMillis(); + } + +} From f4c98595cb973f1c74f76c647d7bf4c4334851e9 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 11:07:59 +0300 Subject: [PATCH 3/9] generator node: comment added on direct executor --- .../org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java index 8fa4b97063..ab9c68e41f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java @@ -152,7 +152,7 @@ public class TbMsgGeneratorNode implements TbNode { ctx.logJsEvalResponse(); prevMsg = ctx.newMsg(ServiceQueue.MAIN, generated.getType(), originatorId, msg.getCustomerId(), generated.getMetaData(), generated.getData()); return Futures.immediateFuture(prevMsg); - }, MoreExecutors.directExecutor()); + }, MoreExecutors.directExecutor()); //usually it runs on js-executor-remote-callback thread pool } return Futures.immediateFuture(prevMsg); From 585f473bdaa9477fd25616982dcc34ba80d50ed9 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 5 May 2021 15:00:00 +0300 Subject: [PATCH 4/9] exception message added onFailure for GeneratorNode --- .../org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java index ab9c68e41f..cb9e25d34b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java @@ -117,7 +117,7 @@ public class TbMsgGeneratorNode implements TbNode { } }, t -> { - log.warn("onMsg onFailure callback, took {}ms, config {}, msg {}", sw.stopAndGetTotalTimeMillis(), config, msg); + 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); scheduleTickMsg(ctx); From 8915b279a027b16c1d45d2e9728ee5dcf474fa7d Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 16 Jun 2021 12:40:38 +0300 Subject: [PATCH 5/9] js-script-engine-api: executeSwitch replaced with asynchronous executeSwitchAsync --- .../controller/RuleChainController.java | 2 +- .../script/RuleNodeJsScriptEngine.java | 12 ++++++--- .../rule/engine/api/ScriptEngine.java | 2 +- .../rule/engine/filter/TbJsSwitchNode.java | 27 ++++++++++++------- .../engine/filter/TbJsSwitchNodeTest.java | 2 +- 5 files changed, 29 insertions(+), 16 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index 3b9af2f0f2..39692922c1 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -400,7 +400,7 @@ public class RuleChainController extends BaseController { output = Boolean.toString(result); break; case "switch": - Set states = engine.executeSwitch(inMsg); + Set states = engine.executeSwitchAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS); output = objectMapper.writeValueAsString(states); break; case "json": diff --git a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java index c688d6978d..e9cf8d621f 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java @@ -195,9 +195,7 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S }, MoreExecutors.directExecutor()); } - @Override - public Set executeSwitch(TbMsg msg) throws ScriptException { - JsonNode result = executeScript(msg); + Set executeSwitchPostProcessFunction(JsonNode result) throws ScriptException { if (result.isTextual()) { return Collections.singleton(result.asText()); } else if (result.isArray()) { @@ -217,6 +215,14 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S } } + @Override + public ListenableFuture> executeSwitchAsync(TbMsg msg) { + log.trace("execute switch async, msg {}", msg); + return Futures.transformAsync(executeScriptAsync(msg), + result -> Futures.immediateFuture(executeSwitchPostProcessFunction(result)), + MoreExecutors.directExecutor()); //usually runs in a callbackExecutor + } + private JsonNode executeScript(TbMsg msg) throws ScriptException { try { String[] inArgs = prepareArgs(msg); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java index a5dcf535e4..da9bd93f2b 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java @@ -35,7 +35,7 @@ public interface ScriptEngine { ListenableFuture executeFilterAsync(TbMsg msg); - Set executeSwitch(TbMsg msg) throws ScriptException; + ListenableFuture> executeSwitchAsync(TbMsg msg); JsonNode executeJson(TbMsg msg) throws ScriptException; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java index f677b90158..5af4ccb29d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java +++ b/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; +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 org.checkerframework.checker.nullness.qual.Nullable; import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.ScriptEngine; @@ -58,17 +62,20 @@ public class TbJsSwitchNode implements TbNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { - ListeningExecutor jsExecutor = ctx.getJsExecutor(); ctx.logJsEvalRequest(); - withCallback(jsExecutor.executeAsync(() -> jsEngine.executeSwitch(msg)), - result -> { - ctx.logJsEvalResponse(); - processSwitch(ctx, msg, result); - }, - t -> { - ctx.logJsEvalFailure(); - ctx.tellFailure(msg, t); - }, ctx.getDbCallbackExecutor()); + Futures.addCallback(jsEngine.executeSwitchAsync(msg), new FutureCallback>() { + @Override + public void onSuccess(@Nullable Set result) { + ctx.logJsEvalResponse(); + processSwitch(ctx, msg, result); + } + + @Override + 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 nextRelations) { diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java index cfacf4cb85..fc2e5aa041 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java @@ -73,7 +73,7 @@ public class TbJsSwitchNodeTest { TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson, ruleChainId, ruleNodeId); mockJsExecutor(); - when(scriptEngine.executeSwitch(msg)).thenReturn(Sets.newHashSet("one", "three")); + when(scriptEngine.executeSwitchAsync(msg)).thenReturn(Futures.immediateFuture(Sets.newHashSet("one", "three"))); node.onMsg(ctx, msg); verify(ctx).getJsExecutor(); From f3757ad127cd0029e347d7b5858a99e4eab62f93 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 16 Jun 2021 18:41:18 +0300 Subject: [PATCH 6/9] js-script-engine-api: refactored executeSwitchAsync --- .../script/RuleNodeJsScriptEngine.java | 24 +++++++++---------- 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java index e9cf8d621f..2e6b48693a 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java @@ -18,7 +18,6 @@ package org.thingsboard.server.service.script; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; 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.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; @@ -32,6 +31,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import javax.script.ScriptException; import java.util.ArrayList; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -145,7 +145,7 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S return Futures.transformAsync(executeScriptAsync(prevMsg), result -> { if (!result.isObject()) { 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 Futures.immediateFuture(unbindMsg(result, prevMsg)); }, MoreExecutors.directExecutor()); @@ -195,31 +195,31 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S }, MoreExecutors.directExecutor()); } - Set executeSwitchPostProcessFunction(JsonNode result) throws ScriptException { + ListenableFuture> executeSwitchPostProcessAsyncFunction(JsonNode result) { if (result.isTextual()) { - return Collections.singleton(result.asText()); - } else if (result.isArray()) { - Set nextStates = Sets.newHashSet(); + return Futures.immediateFuture(Collections.singleton(result.asText())); + } + if (result.isArray()) { + Set nextStates = new HashSet<>(); for (JsonNode val : result) { if (!val.isTextual()) { 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 { nextStates.add(val.asText()); } } - return nextStates; - } else { - log.warn("Wrong result type: {}", result.getNodeType()); - throw new ScriptException("Wrong result type: " + result.getNodeType()); + return Futures.immediateFuture(nextStates); } + log.warn("Wrong result type: {}", result.getNodeType()); + return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + result.getNodeType())); } @Override public ListenableFuture> executeSwitchAsync(TbMsg msg) { log.trace("execute switch async, msg {}", msg); return Futures.transformAsync(executeScriptAsync(msg), - result -> Futures.immediateFuture(executeSwitchPostProcessFunction(result)), + this::executeSwitchPostProcessAsyncFunction, MoreExecutors.directExecutor()); //usually runs in a callbackExecutor } From b397dfb518e101558790127b5e7ad9a7e7aef34f Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 16 Jun 2021 20:49:46 +0300 Subject: [PATCH 7/9] js-script-engine-api: js sync calls replaced with completely async calls (CE API only) --- .../controller/RuleChainController.java | 8 +- .../script/RuleNodeJsScriptEngine.java | 140 +++++++----------- .../rule/engine/api/ScriptEngine.java | 9 +- .../rule/engine/api/TbContext.java | 4 + .../rule/engine/action/TbLogNode.java | 30 ++-- .../rule/engine/filter/TbJsSwitchNode.java | 2 - .../engine/filter/TbJsSwitchNodeTest.java | 17 +-- 7 files changed, 82 insertions(+), 128 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index 39692922c1..fbfd0d595f 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -390,13 +390,13 @@ public class RuleChainController extends BaseController { TbMsg inMsg = TbMsg.newMsg(msgType, null, new TbMsgMetaData(metadata), TbMsgDataType.JSON, data); switch (scriptType) { case "update": - output = msgToOutput(engine.executeUpdate(inMsg)); + output = msgToOutput(engine.executeUpdateAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS)); break; case "generate": output = msgToOutput(engine.executeGenerateAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS)); break; case "filter": - boolean result = engine.executeFilter(inMsg); + boolean result = engine.executeFilterAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS); output = Boolean.toString(result); break; case "switch": @@ -404,11 +404,11 @@ public class RuleChainController extends BaseController { output = objectMapper.writeValueAsString(states); break; case "json": - JsonNode json = engine.executeJson(inMsg); + JsonNode json = engine.executeJsonAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS); output = objectMapper.writeValueAsString(json); break; case "string": - output = engine.executeToString(inMsg); + output = engine.executeToStringAsync(inMsg).get(TIMEOUT, TimeUnit.SECONDS); break; default: throw new IllegalArgumentException("Unsupported script type: " + scriptType); diff --git a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java index 2e6b48693a..17e999a533 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java @@ -106,96 +106,77 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S } } - @Override - public List executeUpdate(TbMsg msg) throws ScriptException { - JsonNode result = executeScript(msg); - if (result.isObject()) { - return Collections.singletonList(unbindMsg(result, msg)); - } else if (result.isArray()){ - List res = new ArrayList<>(result.size()); - result.forEach(jsonObject -> res.add(unbindMsg(jsonObject, msg))); - return res; - } else { - log.warn("Wrong result type: {}", result.getNodeType()); - throw new ScriptException("Wrong result type: " + result.getNodeType()); - } - } - @Override public ListenableFuture> executeUpdateAsync(TbMsg msg) { ListenableFuture result = executeScriptAsync(msg); - return Futures.transformAsync(result, json -> { - if (json.isObject()) { - return Futures.immediateFuture(Collections.singletonList(unbindMsg(json, msg))); - } else if (json.isArray()){ - List 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()); + return Futures.transformAsync(result, + json -> executeUpdateTransform(msg, json), + MoreExecutors.directExecutor()); + } + + ListenableFuture> executeUpdateTransform(TbMsg msg, JsonNode json) { + if (json.isObject()) { + return Futures.immediateFuture(Collections.singletonList(unbindMsg(json, msg))); + } else if (json.isArray()) { + List 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 public ListenableFuture executeGenerateAsync(TbMsg prevMsg) { - log.trace("execute generate async, prevMsg {}", prevMsg); - return Futures.transformAsync(executeScriptAsync(prevMsg), result -> { - if (!result.isObject()) { - log.warn("Wrong result type: {}", result.getNodeType()); - Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + result.getNodeType())); - } - return Futures.immediateFuture(unbindMsg(result, prevMsg)); - }, MoreExecutors.directExecutor()); - + return Futures.transformAsync(executeScriptAsync(prevMsg), + result -> executeGenerateTransform(prevMsg, result), + MoreExecutors.directExecutor()); } - @Override - public JsonNode executeJson(TbMsg msg) throws ScriptException { - return executeScript(msg); + ListenableFuture executeGenerateTransform(TbMsg prevMsg, JsonNode result) { + if (!result.isObject()) { + log.warn("Wrong result type: {}", result.getNodeType()); + Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + result.getNodeType())); + } + return Futures.immediateFuture(unbindMsg(result, prevMsg)); } @Override - public ListenableFuture executeJsonAsync(TbMsg msg) throws ScriptException { + public ListenableFuture executeJsonAsync(TbMsg msg) { return executeScriptAsync(msg); } @Override - public String executeToString(TbMsg msg) throws ScriptException { - JsonNode result = executeScript(msg); - if (!result.isTextual()) { - log.warn("Wrong result type: {}", result.getNodeType()); - throw new ScriptException("Wrong result type: " + result.getNodeType()); - } - return result.asText(); + public ListenableFuture executeToStringAsync(TbMsg msg) { + return Futures.transformAsync(executeScriptAsync(msg), + this::executeToStringTransform, + MoreExecutors.directExecutor()); } - @Override - public boolean executeFilter(TbMsg msg) throws ScriptException { - JsonNode result = executeScript(msg); - if (!result.isBoolean()) { - log.warn("Wrong result type: {}", result.getNodeType()); - throw new ScriptException("Wrong result type: " + result.getNodeType()); + ListenableFuture executeToStringTransform(JsonNode result) { + if (result.isTextual()) { + return Futures.immediateFuture(result.asText()); } - return result.asBoolean(); + log.warn("Wrong result type: {}", result.getNodeType()); + return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + result.getNodeType())); } @Override public ListenableFuture executeFilterAsync(TbMsg msg) { - ListenableFuture result = executeScriptAsync(msg); - return Futures.transformAsync(result, json -> { - if (!json.isBoolean()) { - 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()); + return Futures.transformAsync(executeScriptAsync(msg), + this::executeFilterTransform, + MoreExecutors.directExecutor()); + } + + ListenableFuture executeFilterTransform(JsonNode json) { + if (json.isBoolean()) { + return Futures.immediateFuture(json.asBoolean()); + } + log.warn("Wrong result type: {}", json.getNodeType()); + return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + json.getNodeType())); } - ListenableFuture> executeSwitchPostProcessAsyncFunction(JsonNode result) { + ListenableFuture> executeSwitchTransform(JsonNode result) { if (result.isTextual()) { return Futures.immediateFuture(Collections.singleton(result.asText())); } @@ -217,34 +198,19 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S @Override public ListenableFuture> executeSwitchAsync(TbMsg msg) { - log.trace("execute switch async, msg {}", msg); return Futures.transformAsync(executeScriptAsync(msg), - this::executeSwitchPostProcessAsyncFunction, + this::executeSwitchTransform, MoreExecutors.directExecutor()); //usually runs in a callbackExecutor } - private JsonNode executeScript(TbMsg msg) throws ScriptException { - try { - String[] inArgs = prepareArgs(msg); - String eval = sandboxService.invokeFunction(tenantId, msg.getCustomerId(), this.scriptId, inArgs[0], inArgs[1], inArgs[2]).get().toString(); - return mapper.readTree(eval); - } 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 executeScriptAsync(TbMsg msg) { + ListenableFuture executeScriptAsync(TbMsg msg) { log.trace("execute script async, msg {}", 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 executeScriptAsync(CustomerId customerId, Object... args) { + return Futures.transformAsync(sandboxService.invokeFunction(tenantId, customerId, this.scriptId, args), o -> { try { return Futures.immediateFuture(mapper.readTree(o.toString())); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java index da9bd93f2b..012e135c60 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java @@ -19,14 +19,11 @@ import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.msg.TbMsg; -import javax.script.ScriptException; import java.util.List; import java.util.Set; public interface ScriptEngine { - List executeUpdate(TbMsg msg) throws ScriptException; - ListenableFuture> executeUpdateAsync(TbMsg msg); ListenableFuture executeGenerateAsync(TbMsg prevMsg); @@ -37,11 +34,9 @@ public interface ScriptEngine { ListenableFuture> executeSwitchAsync(TbMsg msg); - JsonNode executeJson(TbMsg msg) throws ScriptException; - - ListenableFuture executeJsonAsync(TbMsg msg) throws ScriptException; + ListenableFuture executeJsonAsync(TbMsg msg); - String executeToString(TbMsg msg) throws ScriptException; + ListenableFuture executeToStringAsync(TbMsg msg); void destroy(); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java index 43e3a5e329..62ed9b0fa0 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java @@ -214,6 +214,10 @@ public interface TbContext { EdgeEventService getEdgeEventService(); + /** + * Js script executors call are completely asynchronous + * */ + @Deprecated ListeningExecutor getJsExecutor(); ListeningExecutor getMailExecutor(); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java index 2b2f0d76a4..c981d5062d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java +++ b/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; +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 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.ScriptEngine; import org.thingsboard.rule.engine.api.TbContext; @@ -55,18 +58,21 @@ public class TbLogNode implements TbNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { - ListeningExecutor jsExecutor = ctx.getJsExecutor(); ctx.logJsEvalRequest(); - withCallback(jsExecutor.executeAsync(() -> jsEngine.executeToString(msg)), - toString -> { - ctx.logJsEvalResponse(); - log.info(toString); - ctx.tellSuccess(msg); - }, - t -> { - ctx.logJsEvalResponse(); - ctx.tellFailure(msg, t); - }); + Futures.addCallback(jsEngine.executeToStringAsync(msg), new FutureCallback() { + @Override + public void onSuccess(@Nullable String result) { + ctx.logJsEvalResponse(); + log.info(result); + ctx.tellSuccess(msg); + } + + @Override + public void onFailure(Throwable t) { + ctx.logJsEvalResponse(); + ctx.tellFailure(msg, t); + } + }, MoreExecutors.directExecutor()); //usually js responses runs on js callback executor } @Override diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java index 5af4ccb29d..bc900a0f2d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java @@ -33,8 +33,6 @@ import org.thingsboard.server.common.msg.TbMsg; import java.util.Set; -import static org.thingsboard.common.util.DonAsynchron.withCallback; - @Slf4j @RuleNode( type = ComponentType.FILTER, diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java index fc2e5aa041..4ac81120b0 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeTest.java +++ b/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()); @Test - public void multipleRoutesAreAllowed() throws TbNodeException, ScriptException { + public void multipleRoutesAreAllowed() throws TbNodeException { initWithScript(); TbMsgMetaData metaData = new TbMsgMetaData(); metaData.putValue("temp", "10"); @@ -72,11 +72,9 @@ public class TbJsSwitchNodeTest { String rawJson = "{\"name\": \"Vit\", \"passed\": 5}"; TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson, ruleChainId, ruleNodeId); - mockJsExecutor(); when(scriptEngine.executeSwitchAsync(msg)).thenReturn(Futures.immediateFuture(Sets.newHashSet("one", "three"))); node.onMsg(ctx, msg); - verify(ctx).getJsExecutor(); verify(ctx).tellNext(msg, Sets.newHashSet("one", "three")); } @@ -92,19 +90,6 @@ public class TbJsSwitchNodeTest { node.init(ctx, nodeConfiguration); } - @SuppressWarnings("unchecked") - private void mockJsExecutor() { - when(ctx.getJsExecutor()).thenReturn(executor); - doAnswer((Answer>>) invocationOnMock -> { - try { - Callable task = (Callable) (invocationOnMock.getArguments())[0]; - return Futures.immediateFuture((Set) 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) { ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); verify(ctx).tellFailure(same(msg), captor.capture()); From 966b0923e7e293c0bd5fb6aa3f98d43aa9dc93de Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 14:12:33 +0300 Subject: [PATCH 8/9] js-script-engine-api: fixed CE merge --- .../server/service/script/RuleNodeJsScriptEngine.java | 1 + .../main/java/org/thingsboard/rule/engine/api/ScriptEngine.java | 2 -- 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java index 17e999a533..4ab81702d5 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java @@ -23,6 +23,7 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; 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.TenantId; import org.thingsboard.server.common.msg.TbMsg; diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java index 012e135c60..7ff63319cd 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/ScriptEngine.java @@ -28,8 +28,6 @@ public interface ScriptEngine { ListenableFuture executeGenerateAsync(TbMsg prevMsg); - boolean executeFilter(TbMsg msg) throws ScriptException; - ListenableFuture executeFilterAsync(TbMsg msg); ListenableFuture> executeSwitchAsync(TbMsg msg); From 29910fad8f6634e4fa7cfc1db2cf74d2e723f7d8 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 13:34:54 +0300 Subject: [PATCH 9/9] TbStopWatch updated licence header for CE --- .../thingsboard/common/util/TbStopWatch.java | 35 ++++++------------- 1 file changed, 10 insertions(+), 25 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java index b7dbd0052c..155ffb5360 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java +++ b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java @@ -1,32 +1,17 @@ /** - * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * Copyright © 2016-2021 The Thingsboard Authors * - * Copyright © 2016-2021 ThingsBoard, Inc. All Rights Reserved. + * 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 * - * NOTICE: All information contained herein is, and remains - * the property of ThingsBoard, Inc. and its suppliers, - * if any. The intellectual and technical concepts contained - * herein are proprietary to ThingsBoard, Inc. - * and its suppliers and may be covered by U.S. and Foreign Patents, - * patents in process, and are protected by trade secret or copyright law. + * http://www.apache.org/licenses/LICENSE-2.0 * - * Dissemination of this information or reproduction of this material is strictly forbidden - * unless prior written permission is obtained from COMPANY. - * - * Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees, - * managers or contractors who have executed Confidentiality and Non-disclosure agreements - * explicitly covering such access. - * - * The copyright notice above does not evidence any actual or intended publication - * or disclosure of this source code, which includes - * information that is confidential and/or proprietary, and is a trade secret, of COMPANY. - * ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE, - * OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT - * THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED, - * AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES. - * THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION - * DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS, - * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. + * 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.common.util;