Browse Source

js-script-engine-api: refactored executeSwitchAsync

pull/4755/head
Sergey Matvienko 5 years ago
parent
commit
f3757ad127
  1. 24
      application/src/main/java/org/thingsboard/server/service/script/RuleNodeJsScriptEngine.java

24
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.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;
@ -32,6 +31,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;
@ -145,7 +145,7 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S
return Futures.transformAsync(executeScriptAsync(prevMsg), result -> { return Futures.transformAsync(executeScriptAsync(prevMsg), result -> {
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 Futures.immediateFuture(unbindMsg(result, prevMsg)); return Futures.immediateFuture(unbindMsg(result, prevMsg));
}, MoreExecutors.directExecutor()); }, MoreExecutors.directExecutor());
@ -195,31 +195,31 @@ public class RuleNodeJsScriptEngine implements org.thingsboard.rule.engine.api.S
}, MoreExecutors.directExecutor()); }, MoreExecutors.directExecutor());
} }
Set<String> executeSwitchPostProcessFunction(JsonNode result) throws ScriptException { ListenableFuture<Set<String>> executeSwitchPostProcessAsyncFunction(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()));
} }
@Override @Override
public ListenableFuture<Set<String>> executeSwitchAsync(TbMsg msg) { public ListenableFuture<Set<String>> executeSwitchAsync(TbMsg msg) {
log.trace("execute switch async, msg {}", msg); log.trace("execute switch async, msg {}", msg);
return Futures.transformAsync(executeScriptAsync(msg), return Futures.transformAsync(executeScriptAsync(msg),
result -> Futures.immediateFuture(executeSwitchPostProcessFunction(result)), this::executeSwitchPostProcessAsyncFunction,
MoreExecutors.directExecutor()); //usually runs in a callbackExecutor MoreExecutors.directExecutor()); //usually runs in a callbackExecutor
} }

Loading…
Cancel
Save