198 changed files with 4665 additions and 1297 deletions
@ -1,249 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.service.script; |
|
||||
|
|
||||
import com.google.common.hash.Hashing; |
|
||||
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.Value; |
|
||||
import org.springframework.data.util.Pair; |
|
||||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
||||
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
|
||||
import org.thingsboard.server.common.data.id.CustomerId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|
||||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
||||
|
|
||||
import java.util.Map; |
|
||||
import java.util.UUID; |
|
||||
import java.util.concurrent.ConcurrentHashMap; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.ScheduledExecutorService; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
|
|
||||
import static java.lang.String.format; |
|
||||
|
|
||||
/** |
|
||||
* Created by ashvayka on 26.09.18. |
|
||||
*/ |
|
||||
@Slf4j |
|
||||
@SuppressWarnings("UnstableApiUsage") |
|
||||
public abstract class AbstractJsInvokeService implements JsInvokeService { |
|
||||
|
|
||||
private final TbApiUsageStateService apiUsageStateService; |
|
||||
private final TbApiUsageClient apiUsageClient; |
|
||||
protected ScheduledExecutorService timeoutExecutorService; |
|
||||
|
|
||||
protected final Map<UUID, Pair<String, String>> scriptIdToNameAndHashMap = new ConcurrentHashMap<>(); |
|
||||
protected final Map<UUID, DisableListInfo> disabledFunctions = new ConcurrentHashMap<>(); |
|
||||
|
|
||||
@Getter |
|
||||
@Value("${js.max_total_args_size:100000}") |
|
||||
private long maxTotalArgsSize; |
|
||||
@Getter |
|
||||
@Value("${js.max_result_size:300000}") |
|
||||
private long maxResultSize; |
|
||||
@Getter |
|
||||
@Value("${js.max_script_body_size:50000}") |
|
||||
private long maxScriptBodySize; |
|
||||
|
|
||||
protected AbstractJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient) { |
|
||||
this.apiUsageStateService = apiUsageStateService; |
|
||||
this.apiUsageClient = apiUsageClient; |
|
||||
} |
|
||||
|
|
||||
public void init(long maxRequestsTimeout) { |
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
timeoutExecutorService = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("nashorn-js-timeout")); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void stop() { |
|
||||
if (timeoutExecutorService != null) { |
|
||||
timeoutExecutorService.shutdownNow(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<UUID> eval(TenantId tenantId, JsScriptType scriptType, String scriptBody, String... argNames) { |
|
||||
if (apiUsageStateService.getApiUsageState(tenantId).isJsExecEnabled()) { |
|
||||
if (scriptBodySizeExceeded(scriptBody)) { |
|
||||
return error(format("Script body exceeds maximum allowed size of %s symbols", getMaxScriptBodySize())); |
|
||||
} |
|
||||
UUID scriptId = UUID.randomUUID(); |
|
||||
String scriptHash = hash(tenantId, scriptBody); |
|
||||
String functionName = constructFunctionName(scriptId, scriptHash); |
|
||||
String jsScript = generateJsScript(scriptType, functionName, scriptBody, argNames); |
|
||||
return doEval(scriptId, scriptHash, functionName, jsScript); |
|
||||
} else { |
|
||||
return error("JS Execution is disabled due to API limits!"); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
protected String constructFunctionName(UUID scriptId, String scriptHash) { |
|
||||
return "invokeInternal_" + scriptId.toString().replace('-', '_'); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<String> invokeFunction(TenantId tenantId, CustomerId customerId, UUID scriptId, Object... args) { |
|
||||
if (apiUsageStateService.getApiUsageState(tenantId).isJsExecEnabled()) { |
|
||||
Pair<String, String> nameAndHash = scriptIdToNameAndHashMap.get(scriptId); |
|
||||
if (nameAndHash == null) { |
|
||||
return error("No compiled script found for scriptId: [" + scriptId + "]!"); |
|
||||
} |
|
||||
String functionName = nameAndHash.getFirst(); |
|
||||
String scriptHash = nameAndHash.getSecond(); |
|
||||
if (!isDisabled(scriptId)) { |
|
||||
if (argsSizeExceeded(args)) { |
|
||||
return scriptExecutionError(scriptId, format("Script input arguments exceed maximum allowed total args size of %s symbols", getMaxTotalArgsSize())); |
|
||||
} |
|
||||
apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.JS_EXEC_COUNT, 1); |
|
||||
return Futures.transformAsync(doInvokeFunction(scriptId, scriptHash, functionName, args), output -> { |
|
||||
String result = output.toString(); |
|
||||
if (resultSizeExceeded(result)) { |
|
||||
return scriptExecutionError(scriptId, format("Script invocation result exceeds maximum allowed size of %s symbols", getMaxResultSize())); |
|
||||
} |
|
||||
return Futures.immediateFuture(result); |
|
||||
}, MoreExecutors.directExecutor()); |
|
||||
} else { |
|
||||
String message = "Script invocation is blocked due to maximum error count " |
|
||||
+ getMaxErrors() + ", scriptId " + scriptId + "!"; |
|
||||
log.warn(message); |
|
||||
return error(message); |
|
||||
} |
|
||||
} else { |
|
||||
return error("JS Execution is disabled due to API limits!"); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Void> release(UUID scriptId) { |
|
||||
Pair<String, String> nameAndHash = scriptIdToNameAndHashMap.get(scriptId); |
|
||||
if (nameAndHash != null) { |
|
||||
try { |
|
||||
scriptIdToNameAndHashMap.remove(scriptId); |
|
||||
disabledFunctions.remove(scriptId); |
|
||||
doRelease(scriptId, nameAndHash.getSecond(), nameAndHash.getFirst()); |
|
||||
} catch (Exception e) { |
|
||||
return Futures.immediateFailedFuture(e); |
|
||||
} |
|
||||
} |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
|
|
||||
protected abstract ListenableFuture<UUID> doEval(UUID scriptId, String scriptHash, String functionName, String scriptBody); |
|
||||
|
|
||||
protected abstract ListenableFuture<Object> doInvokeFunction(UUID scriptId, String scriptHash, String functionName, Object[] args); |
|
||||
|
|
||||
protected abstract void doRelease(UUID scriptId, String scriptHash, String functionName) throws Exception; |
|
||||
|
|
||||
protected abstract int getMaxErrors(); |
|
||||
|
|
||||
protected abstract long getMaxBlacklistDuration(); |
|
||||
|
|
||||
protected String hash(TenantId tenantId, String scriptBody) { |
|
||||
return Hashing.murmur3_128().newHasher() |
|
||||
.putLong(tenantId.getId().getMostSignificantBits()) |
|
||||
.putLong(tenantId.getId().getLeastSignificantBits()) |
|
||||
.putUnencodedChars(scriptBody) |
|
||||
.hash().toString(); |
|
||||
} |
|
||||
|
|
||||
protected void onScriptExecutionError(UUID scriptId, Throwable t, String scriptBody) { |
|
||||
DisableListInfo disableListInfo = disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()); |
|
||||
log.warn("Script has exception and will increment counter {} on disabledFunctions for id {}, exception {}, cause {}, scriptBody {}", |
|
||||
disableListInfo.get(), scriptId, t, t.getCause(), scriptBody); |
|
||||
disableListInfo.incrementAndGet(); |
|
||||
} |
|
||||
|
|
||||
private boolean scriptBodySizeExceeded(String scriptBody) { |
|
||||
if (getMaxScriptBodySize() <= 0) return false; |
|
||||
return scriptBody.length() > getMaxScriptBodySize(); |
|
||||
} |
|
||||
|
|
||||
private boolean argsSizeExceeded(Object[] args) { |
|
||||
if (getMaxTotalArgsSize() <= 0) return false; |
|
||||
long totalArgsSize = 0; |
|
||||
for (Object arg : args) { |
|
||||
if (arg instanceof CharSequence) { |
|
||||
totalArgsSize += ((CharSequence) arg).length(); |
|
||||
} |
|
||||
} |
|
||||
return totalArgsSize > getMaxTotalArgsSize(); |
|
||||
} |
|
||||
|
|
||||
private boolean resultSizeExceeded(String result) { |
|
||||
if (getMaxResultSize() <= 0) return false; |
|
||||
return result.length() > getMaxResultSize(); |
|
||||
} |
|
||||
|
|
||||
private String generateJsScript(JsScriptType scriptType, String functionName, String scriptBody, String... argNames) { |
|
||||
if (scriptType == JsScriptType.RULE_NODE_SCRIPT) { |
|
||||
return RuleNodeScriptFactory.generateRuleNodeScript(functionName, scriptBody, argNames); |
|
||||
} |
|
||||
throw new RuntimeException("No script factory implemented for scriptType: " + scriptType); |
|
||||
} |
|
||||
|
|
||||
private boolean isDisabled(UUID scriptId) { |
|
||||
DisableListInfo errorCount = disabledFunctions.get(scriptId); |
|
||||
if (errorCount != null) { |
|
||||
if (errorCount.getExpirationTime() <= System.currentTimeMillis()) { |
|
||||
disabledFunctions.remove(scriptId); |
|
||||
return false; |
|
||||
} else { |
|
||||
return errorCount.get() >= getMaxErrors(); |
|
||||
} |
|
||||
} else { |
|
||||
return false; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private <T> ListenableFuture<T> error(String message) { |
|
||||
return Futures.immediateFailedFuture(new RuntimeException(message)); |
|
||||
} |
|
||||
|
|
||||
private <T> ListenableFuture<T> scriptExecutionError(UUID scriptId, String errorMsg) { |
|
||||
RuntimeException error = new RuntimeException(errorMsg); |
|
||||
onScriptExecutionError(scriptId, error, null); |
|
||||
return Futures.immediateFailedFuture(error); |
|
||||
} |
|
||||
|
|
||||
private class DisableListInfo { |
|
||||
private final AtomicInteger counter; |
|
||||
private long expirationTime; |
|
||||
|
|
||||
private DisableListInfo() { |
|
||||
this.counter = new AtomicInteger(0); |
|
||||
} |
|
||||
|
|
||||
public int get() { |
|
||||
return counter.get(); |
|
||||
} |
|
||||
|
|
||||
public int incrementAndGet() { |
|
||||
int result = counter.incrementAndGet(); |
|
||||
expirationTime = System.currentTimeMillis() + getMaxBlacklistDuration(); |
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
public long getExpirationTime() { |
|
||||
return expirationTime; |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,185 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.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 delight.nashornsandbox.NashornSandbox; |
|
||||
import delight.nashornsandbox.NashornSandboxes; |
|
||||
import lombok.Getter; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.data.util.Pair; |
|
||||
import org.springframework.scheduling.annotation.Scheduled; |
|
||||
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
||||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|
||||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
||||
|
|
||||
import javax.annotation.PostConstruct; |
|
||||
import javax.annotation.PreDestroy; |
|
||||
import javax.script.Invocable; |
|
||||
import javax.script.ScriptEngine; |
|
||||
import javax.script.ScriptEngineManager; |
|
||||
import javax.script.ScriptException; |
|
||||
import java.util.UUID; |
|
||||
import java.util.concurrent.ExecutionException; |
|
||||
import java.util.concurrent.ExecutorService; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
import java.util.concurrent.locks.ReentrantLock; |
|
||||
|
|
||||
@Slf4j |
|
||||
public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeService { |
|
||||
|
|
||||
private NashornSandbox sandbox; |
|
||||
private ScriptEngine engine; |
|
||||
private ExecutorService monitorExecutorService; |
|
||||
|
|
||||
private final AtomicInteger jsPushedMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsInvokeMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsEvalMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsFailedMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsTimeoutMsgs = new AtomicInteger(0); |
|
||||
private final FutureCallback<UUID> evalCallback = new JsStatCallback<>(jsEvalMsgs, jsTimeoutMsgs, jsFailedMsgs); |
|
||||
private final FutureCallback<Object> invokeCallback = new JsStatCallback<>(jsInvokeMsgs, jsTimeoutMsgs, jsFailedMsgs); |
|
||||
|
|
||||
private final ReentrantLock evalLock = new ReentrantLock(); |
|
||||
|
|
||||
@Getter |
|
||||
private final JsExecutorService jsExecutor; |
|
||||
|
|
||||
@Value("${js.local.max_requests_timeout:0}") |
|
||||
private long maxRequestsTimeout; |
|
||||
|
|
||||
@Value("${js.local.stats.enabled:false}") |
|
||||
private boolean statsEnabled; |
|
||||
|
|
||||
public AbstractNashornJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient, JsExecutorService jsExecutor) { |
|
||||
super(apiUsageStateService, apiUsageClient); |
|
||||
this.jsExecutor = jsExecutor; |
|
||||
} |
|
||||
|
|
||||
@Scheduled(fixedDelayString = "${js.local.stats.print_interval_ms:10000}") |
|
||||
public void printStats() { |
|
||||
if (statsEnabled) { |
|
||||
int pushedMsgs = jsPushedMsgs.getAndSet(0); |
|
||||
int invokeMsgs = jsInvokeMsgs.getAndSet(0); |
|
||||
int evalMsgs = jsEvalMsgs.getAndSet(0); |
|
||||
int failed = jsFailedMsgs.getAndSet(0); |
|
||||
int timedOut = jsTimeoutMsgs.getAndSet(0); |
|
||||
if (pushedMsgs > 0 || invokeMsgs > 0 || evalMsgs > 0 || failed > 0 || timedOut > 0) { |
|
||||
log.info("Nashorn JS Invoke Stats: pushed [{}] received [{}] invoke [{}] eval [{}] failed [{}] timedOut [{}]", |
|
||||
pushedMsgs, invokeMsgs + evalMsgs, invokeMsgs, evalMsgs, failed, timedOut); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PostConstruct |
|
||||
public void init() { |
|
||||
super.init(maxRequestsTimeout); |
|
||||
if (useJsSandbox()) { |
|
||||
sandbox = NashornSandboxes.create(); |
|
||||
monitorExecutorService = ThingsBoardExecutors.newWorkStealingPool(getMonitorThreadPoolSize(), "nashorn-js-monitor"); |
|
||||
sandbox.setExecutor(monitorExecutorService); |
|
||||
sandbox.setMaxCPUTime(getMaxCpuTime()); |
|
||||
sandbox.allowNoBraces(false); |
|
||||
sandbox.allowLoadFunctions(true); |
|
||||
sandbox.setMaxPreparedStatements(30); |
|
||||
} else { |
|
||||
ScriptEngineManager factory = new ScriptEngineManager(); |
|
||||
engine = factory.getEngineByName("nashorn"); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreDestroy |
|
||||
public void stop() { |
|
||||
super.stop(); |
|
||||
if (monitorExecutorService != null) { |
|
||||
monitorExecutorService.shutdownNow(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
protected abstract boolean useJsSandbox(); |
|
||||
|
|
||||
protected abstract int getMonitorThreadPoolSize(); |
|
||||
|
|
||||
protected abstract long getMaxCpuTime(); |
|
||||
|
|
||||
@Override |
|
||||
protected ListenableFuture<UUID> doEval(UUID scriptId, String scriptHash, String functionName, String jsScript) { |
|
||||
jsPushedMsgs.incrementAndGet(); |
|
||||
ListenableFuture<UUID> result = jsExecutor.executeAsync(() -> { |
|
||||
try { |
|
||||
evalLock.lock(); |
|
||||
try { |
|
||||
if (useJsSandbox()) { |
|
||||
sandbox.eval(jsScript); |
|
||||
} else { |
|
||||
engine.eval(jsScript); |
|
||||
} |
|
||||
} finally { |
|
||||
evalLock.unlock(); |
|
||||
} |
|
||||
scriptIdToNameAndHashMap.put(scriptId, Pair.of(functionName, scriptHash)); |
|
||||
return scriptId; |
|
||||
} catch (Exception e) { |
|
||||
log.debug("Failed to compile JS script: {}", e.getMessage(), e); |
|
||||
throw new ExecutionException(e); |
|
||||
} |
|
||||
}); |
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
result = Futures.withTimeout(result, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
||||
} |
|
||||
Futures.addCallback(result, evalCallback, MoreExecutors.directExecutor()); |
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected ListenableFuture<Object> doInvokeFunction(UUID scriptId, String scriptHash, String functionName, Object[] args) { |
|
||||
jsPushedMsgs.incrementAndGet(); |
|
||||
ListenableFuture<Object> result = jsExecutor.executeAsync(() -> { |
|
||||
try { |
|
||||
if (useJsSandbox()) { |
|
||||
return sandbox.getSandboxedInvocable().invokeFunction(functionName, args); |
|
||||
} else { |
|
||||
return ((Invocable) engine).invokeFunction(functionName, args); |
|
||||
} |
|
||||
} catch (ScriptException e) { |
|
||||
throw new ExecutionException(e); |
|
||||
} catch (Exception e) { |
|
||||
onScriptExecutionError(scriptId, e, functionName); |
|
||||
throw new ExecutionException(e); |
|
||||
} |
|
||||
}); |
|
||||
|
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
result = Futures.withTimeout(result, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
||||
} |
|
||||
Futures.addCallback(result, invokeCallback, MoreExecutors.directExecutor()); |
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
protected void doRelease(UUID scriptId, String scriptHash, String functionName) throws ScriptException { |
|
||||
if (useJsSandbox()) { |
|
||||
sandbox.eval(functionName + " = undefined;"); |
|
||||
} else { |
|
||||
engine.eval(functionName + " = undefined;"); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,64 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.service.script; |
|
||||
|
|
||||
import lombok.Getter; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
|
||||
import org.springframework.stereotype.Service; |
|
||||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|
||||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
||||
|
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
|
|
||||
@Slf4j |
|
||||
@ConditionalOnProperty(prefix = "js", value = "evaluator", havingValue = "local", matchIfMissing = true) |
|
||||
@Service |
|
||||
public class NashornJsInvokeService extends AbstractNashornJsInvokeService { |
|
||||
|
|
||||
@Value("${js.local.use_js_sandbox}") |
|
||||
private boolean useJsSandbox; |
|
||||
|
|
||||
@Getter |
|
||||
@Value("${js.local.monitor_thread_pool_size}") |
|
||||
private int monitorThreadPoolSize; |
|
||||
|
|
||||
@Getter |
|
||||
@Value("${js.local.max_cpu_time}") |
|
||||
private long maxCpuTime; |
|
||||
|
|
||||
@Getter |
|
||||
@Value("${js.local.max_errors}") |
|
||||
private int maxErrors; |
|
||||
|
|
||||
@Value("${js.local.max_black_list_duration_sec:60}") |
|
||||
private int maxBlackListDurationSec; |
|
||||
|
|
||||
public NashornJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient, JsExecutorService jsExecutor) { |
|
||||
super(apiUsageStateService, apiUsageClient, jsExecutor); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected boolean useJsSandbox() { |
|
||||
return useJsSandbox; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected long getMaxBlacklistDuration() { |
|
||||
return TimeUnit.SECONDS.toMillis(maxBlackListDurationSec); |
|
||||
} |
|
||||
} |
|
||||
@ -0,0 +1,171 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.service.script; |
||||
|
|
||||
|
import com.fasterxml.jackson.core.type.TypeReference; |
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
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.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.script.api.RuleNodeScriptFactory; |
||||
|
import org.thingsboard.script.api.mvel.MvelInvokeService; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
||||
|
|
||||
|
import javax.script.ScriptException; |
||||
|
import java.util.ArrayList; |
||||
|
import java.util.Collection; |
||||
|
import java.util.Collections; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.HashSet; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.Set; |
||||
|
import java.util.stream.Collectors; |
||||
|
|
||||
|
|
||||
|
@Slf4j |
||||
|
public class RuleNodeMvelScriptEngine extends RuleNodeScriptEngine<MvelInvokeService, Object> { |
||||
|
|
||||
|
public RuleNodeMvelScriptEngine(TenantId tenantId, MvelInvokeService scriptInvokeService, String script, String... argNames) { |
||||
|
super(tenantId, scriptInvokeService, script, argNames); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<Boolean> executeFilterTransform(Object result) { |
||||
|
if (result instanceof Boolean) { |
||||
|
return Futures.immediateFuture((Boolean) result); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<List<TbMsg>> executeUpdateTransform(TbMsg msg, Object result) { |
||||
|
if (result instanceof Map) { |
||||
|
return Futures.immediateFuture(Collections.singletonList(unbindMsg((Map) result, msg))); |
||||
|
} else if (result instanceof Collection) { |
||||
|
List<TbMsg> res = new ArrayList<>(); |
||||
|
for (Object resObject : (Collection) result) { |
||||
|
if (resObject instanceof Map) { |
||||
|
res.add(unbindMsg((Map) result, msg)); |
||||
|
} else { |
||||
|
return wrongResultType(resObject); |
||||
|
} |
||||
|
} |
||||
|
return Futures.immediateFuture(res); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<TbMsg> executeGenerateTransform(TbMsg prevMsg, Object result) { |
||||
|
if (result instanceof Map) { |
||||
|
return Futures.immediateFuture(unbindMsg((Map) result, prevMsg)); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<String> executeToStringTransform(Object result) { |
||||
|
if (result instanceof String) { |
||||
|
return Futures.immediateFuture((String) result); |
||||
|
} else { |
||||
|
return Futures.immediateFuture(JacksonUtil.toString(result)); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<Set<String>> executeSwitchTransform(Object result) { |
||||
|
if (result instanceof String) { |
||||
|
return Futures.immediateFuture(Collections.singleton((String) result)); |
||||
|
} else if (result instanceof Collection) { |
||||
|
Set<String> res = new HashSet<>(); |
||||
|
for (Object resObject : (Collection) result) { |
||||
|
if (resObject instanceof String) { |
||||
|
res.add((String) resObject); |
||||
|
} else { |
||||
|
return wrongResultType(resObject); |
||||
|
} |
||||
|
} |
||||
|
return Futures.immediateFuture(res); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<JsonNode> executeJsonAsync(TbMsg msg) { |
||||
|
return Futures.transform(executeScriptAsync(msg), JacksonUtil::valueToTree, MoreExecutors.directExecutor()); |
||||
|
|
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected Object convertResult(Object result) { |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected Object[] prepareArgs(TbMsg msg) { |
||||
|
Object[] args = new Object[3]; |
||||
|
if (msg.getData() != null) { |
||||
|
args[0] = JacksonUtil.fromString(msg.getData(), Map.class); |
||||
|
} else { |
||||
|
args[0] = new HashMap<>(); |
||||
|
} |
||||
|
args[1] = new HashMap<>(msg.getMetaData().getData()); |
||||
|
args[2] = msg.getType(); |
||||
|
return args; |
||||
|
} |
||||
|
|
||||
|
private static TbMsg unbindMsg(Map msgData, TbMsg msg) { |
||||
|
String data = null; |
||||
|
Map<String, String> metadata = null; |
||||
|
String messageType = null; |
||||
|
if (msgData.containsKey(RuleNodeScriptFactory.MSG)) { |
||||
|
data = JacksonUtil.toString(msgData.get(RuleNodeScriptFactory.MSG)); |
||||
|
} |
||||
|
if (msgData.containsKey(RuleNodeScriptFactory.METADATA)) { |
||||
|
Object msgMetadataObj = msgData.get(RuleNodeScriptFactory.METADATA); |
||||
|
if (msgMetadataObj instanceof Map) { |
||||
|
metadata = ((Map<?, ?>) msgMetadataObj).entrySet().stream().filter(e -> e.getValue() != null) |
||||
|
.collect(Collectors.toMap(e -> e.getKey().toString(), e -> e.getValue().toString())); |
||||
|
} else { |
||||
|
metadata = JacksonUtil.convertValue(msgMetadataObj, new TypeReference<>() { |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
|
if (msgData.containsKey(RuleNodeScriptFactory.MSG_TYPE)) { |
||||
|
messageType = msgData.get(RuleNodeScriptFactory.MSG_TYPE).toString(); |
||||
|
} |
||||
|
String newData = data != null ? data : msg.getData(); |
||||
|
TbMsgMetaData newMetadata = metadata != null ? new TbMsgMetaData(metadata) : msg.getMetaData().copy(); |
||||
|
String newMessageType = !StringUtils.isEmpty(messageType) ? messageType : msg.getType(); |
||||
|
return TbMsg.transformMsg(msg, newMessageType, msg.getOriginator(), newMetadata, newData); |
||||
|
} |
||||
|
|
||||
|
private static <T> ListenableFuture<T> wrongResultType(Object result) { |
||||
|
String className = toClassName(result); |
||||
|
log.warn("Wrong result type: {}", className); |
||||
|
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + className)); |
||||
|
} |
||||
|
|
||||
|
private static String toClassName(Object result) { |
||||
|
return result != null ? result.getClass().getSimpleName() : "null"; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,133 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.service.script; |
||||
|
|
||||
|
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.thingsboard.rule.engine.api.ScriptEngine; |
||||
|
import org.thingsboard.script.api.ScriptInvokeService; |
||||
|
import org.thingsboard.script.api.ScriptType; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
|
||||
|
import javax.script.ScriptException; |
||||
|
import java.util.List; |
||||
|
import java.util.Set; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
|
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class RuleNodeScriptEngine<T extends ScriptInvokeService, R> implements ScriptEngine { |
||||
|
|
||||
|
private final T scriptInvokeService; |
||||
|
|
||||
|
private final UUID scriptId; |
||||
|
private final TenantId tenantId; |
||||
|
|
||||
|
public RuleNodeScriptEngine(TenantId tenantId, T scriptInvokeService, String script, String... argNames) { |
||||
|
this.tenantId = tenantId; |
||||
|
this.scriptInvokeService = scriptInvokeService; |
||||
|
try { |
||||
|
this.scriptId = this.scriptInvokeService.eval(tenantId, ScriptType.RULE_NODE_SCRIPT, script, argNames).get(); |
||||
|
} catch (Exception e) { |
||||
|
Throwable t = e; |
||||
|
if (e instanceof ExecutionException) { |
||||
|
t = e.getCause(); |
||||
|
} |
||||
|
throw new IllegalArgumentException("Can't compile script: " + t.getMessage(), t); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected abstract Object[] prepareArgs(TbMsg msg); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<List<TbMsg>> executeUpdateAsync(TbMsg msg) { |
||||
|
ListenableFuture<R> result = executeScriptAsync(msg); |
||||
|
return Futures.transformAsync(result, |
||||
|
json -> executeUpdateTransform(msg, json), |
||||
|
MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected abstract ListenableFuture<List<TbMsg>> executeUpdateTransform(TbMsg msg, R result); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<TbMsg> executeGenerateAsync(TbMsg prevMsg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(prevMsg), |
||||
|
result -> executeGenerateTransform(prevMsg, result), |
||||
|
MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected abstract ListenableFuture<TbMsg> executeGenerateTransform(TbMsg prevMsg, R result); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<String> executeToStringAsync(TbMsg msg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(msg), this::executeToStringTransform, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Boolean> executeFilterAsync(TbMsg msg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(msg), |
||||
|
this::executeFilterTransform, |
||||
|
MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected abstract ListenableFuture<String> executeToStringTransform(R result); |
||||
|
|
||||
|
protected abstract ListenableFuture<Boolean> executeFilterTransform(R result); |
||||
|
|
||||
|
protected abstract ListenableFuture<Set<String>> executeSwitchTransform(R result); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Set<String>> executeSwitchAsync(TbMsg msg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(msg), |
||||
|
this::executeSwitchTransform, |
||||
|
MoreExecutors.directExecutor()); //usually runs in a callbackExecutor
|
||||
|
} |
||||
|
|
||||
|
ListenableFuture<R> executeScriptAsync(TbMsg msg) { |
||||
|
log.trace("execute script async, msg {}", msg); |
||||
|
Object[] inArgs = prepareArgs(msg); |
||||
|
return executeScriptAsync(msg.getCustomerId(), inArgs[0], inArgs[1], inArgs[2]); |
||||
|
} |
||||
|
|
||||
|
ListenableFuture<R> executeScriptAsync(CustomerId customerId, Object... args) { |
||||
|
return Futures.transformAsync(scriptInvokeService.invokeScript(tenantId, customerId, this.scriptId, args), |
||||
|
o -> { |
||||
|
try { |
||||
|
return Futures.immediateFuture(convertResult(o)); |
||||
|
} catch (Exception e) { |
||||
|
if (e.getCause() instanceof ScriptException) { |
||||
|
return Futures.immediateFailedFuture(e.getCause()); |
||||
|
} else if (e.getCause() instanceof RuntimeException) { |
||||
|
return Futures.immediateFailedFuture(new ScriptException(e.getCause().getMessage())); |
||||
|
} else { |
||||
|
return Futures.immediateFailedFuture(new ScriptException(e)); |
||||
|
} |
||||
|
} |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
public void destroy() { |
||||
|
scriptInvokeService.release(this.scriptId); |
||||
|
} |
||||
|
|
||||
|
protected abstract R convertResult(Object result); |
||||
|
} |
||||
@ -0,0 +1,128 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.service.script; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.node.ObjectNode; |
||||
|
import org.junit.Assert; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.test.context.TestPropertySource; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.script.api.ScriptType; |
||||
|
import org.thingsboard.script.api.mvel.MvelInvokeService; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.controller.AbstractControllerTest; |
||||
|
import org.thingsboard.server.dao.service.DaoSqlTest; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
||||
|
|
||||
|
@DaoSqlTest |
||||
|
@TestPropertySource(properties = { |
||||
|
"mvel.max_script_body_size=100", |
||||
|
"mvel.max_total_args_size=50", |
||||
|
"mvel.max_result_size=50", |
||||
|
"mvel.max_errors=2", |
||||
|
}) |
||||
|
class MvelInvokeServiceTest extends AbstractControllerTest { |
||||
|
|
||||
|
@Autowired |
||||
|
private MvelInvokeService invokeService; |
||||
|
|
||||
|
@Value("${mvel.max_errors}") |
||||
|
private int maxJsErrors; |
||||
|
|
||||
|
@Test |
||||
|
void givenSimpleScriptTestPerformance() throws ExecutionException, InterruptedException { |
||||
|
int iterations = 100000; |
||||
|
UUID scriptId = evalScript("return msg.temperature > 20"); |
||||
|
// warmup
|
||||
|
ObjectNode msg = JacksonUtil.newObjectNode(); |
||||
|
for (int i = 0; i < 100; i++) { |
||||
|
msg.put("temperature", i); |
||||
|
boolean expected = i > 20; |
||||
|
boolean result = Boolean.valueOf(invokeScript(scriptId, JacksonUtil.toString(msg))); |
||||
|
Assert.assertEquals(expected, result); |
||||
|
} |
||||
|
long startTs = System.currentTimeMillis(); |
||||
|
for (int i = 0; i < iterations; i++) { |
||||
|
msg.put("temperature", i); |
||||
|
boolean expected = i > 20; |
||||
|
boolean result = Boolean.valueOf(invokeScript(scriptId, JacksonUtil.toString(msg))); |
||||
|
Assert.assertEquals(expected, result); |
||||
|
} |
||||
|
long duration = System.currentTimeMillis() - startTs; |
||||
|
System.out.println(iterations + " invocations took: " + duration + "ms"); |
||||
|
Assert.assertTrue(duration < TimeUnit.MINUTES.toMillis(1)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenTooBigScriptForEval_thenReturnError() { |
||||
|
String hugeScript = "var a = 'qwertyqwertywertyqwabababerqwertyqwertywertyqwabababerqwertyqwertywertyqwabababerqwertyqwertywertyqwabababerqwertyqwertywertyqwabababer'; return {a: a};"; |
||||
|
|
||||
|
assertThatThrownBy(() -> { |
||||
|
evalScript(hugeScript); |
||||
|
}).hasMessageContaining("body exceeds maximum allowed size"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenTooBigScriptInputArgs_thenReturnErrorAndReportScriptExecutionError() throws Exception { |
||||
|
String script = "return { msg: msg };"; |
||||
|
String hugeMsg = "{\"input\":\"123456781234349\"}"; |
||||
|
UUID scriptId = evalScript(script); |
||||
|
|
||||
|
for (int i = 0; i < maxJsErrors; i++) { |
||||
|
assertThatThrownBy(() -> { |
||||
|
invokeScript(scriptId, hugeMsg); |
||||
|
}).hasMessageContaining("input arguments exceed maximum"); |
||||
|
} |
||||
|
assertThatScriptIsBlocked(scriptId); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void whenScriptInvocationResultIsTooBig_thenReturnErrorAndReportScriptExecutionError() throws Exception { |
||||
|
String script = "var s = 'a'; for(int i=0; i<50; i++){ s +='a';} return { s: s};"; |
||||
|
UUID scriptId = evalScript(script); |
||||
|
|
||||
|
for (int i = 0; i < maxJsErrors; i++) { |
||||
|
assertThatThrownBy(() -> { |
||||
|
invokeScript(scriptId, "{}"); |
||||
|
}).hasMessageContaining("result exceeds maximum allowed size"); |
||||
|
} |
||||
|
assertThatScriptIsBlocked(scriptId); |
||||
|
} |
||||
|
|
||||
|
private void assertThatScriptIsBlocked(UUID scriptId) { |
||||
|
assertThatThrownBy(() -> { |
||||
|
invokeScript(scriptId, "{}"); |
||||
|
}).hasMessageContaining("invocation is blocked due to maximum error"); |
||||
|
} |
||||
|
|
||||
|
private UUID evalScript(String script) throws ExecutionException, InterruptedException { |
||||
|
return invokeService.eval(TenantId.SYS_TENANT_ID, ScriptType.RULE_NODE_SCRIPT, script, "msg", "metadata", "msgType").get(); |
||||
|
} |
||||
|
|
||||
|
private String invokeScript(UUID scriptId, String str) throws ExecutionException, InterruptedException { |
||||
|
var msg = JacksonUtil.fromString(str, Map.class); |
||||
|
return invokeService.invokeScript(TenantId.SYS_TENANT_ID, null, scriptId, msg, "{}", "POST_TELEMETRY_REQUEST").get().toString(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,196 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.system; |
||||
|
|
||||
|
import com.fasterxml.jackson.core.type.TypeReference; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.ListeningExecutorService; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import com.google.common.util.concurrent.Uninterruptibles; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.junit.After; |
||||
|
import org.junit.Assert; |
||||
|
import org.junit.Before; |
||||
|
import org.junit.Test; |
||||
|
import org.springframework.test.web.servlet.ResultActions; |
||||
|
import org.thingsboard.server.common.data.TenantProfile; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; |
||||
|
import org.thingsboard.server.common.data.tenant.profile.TenantProfileConfiguration; |
||||
|
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; |
||||
|
import org.thingsboard.server.controller.AbstractControllerTest; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
import java.util.Optional; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
import java.util.concurrent.Executors; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.concurrent.TimeoutException; |
||||
|
|
||||
|
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; |
||||
|
|
||||
|
/** |
||||
|
* @author Illia Barkov |
||||
|
*/ |
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class BaseRestApiLimitsTest extends AbstractControllerTest { |
||||
|
|
||||
|
private static int MESSAGES_LIMIT = 10; |
||||
|
private static int TIME_FOR_LIMIT = 5; |
||||
|
|
||||
|
TenantProfile tenantProfile; |
||||
|
|
||||
|
ListeningExecutorService service = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(10)); |
||||
|
|
||||
|
@Before |
||||
|
public void before() throws Exception { |
||||
|
loginSysAdmin(); |
||||
|
tenantProfile = getDefaultTenantProfile(); |
||||
|
logout(); |
||||
|
} |
||||
|
|
||||
|
@After |
||||
|
public void after() throws Exception { |
||||
|
logout(); |
||||
|
loginSysAdmin(); |
||||
|
saveTenantProfileWitConfiguration(tenantProfile, new DefaultTenantProfileConfiguration()); |
||||
|
logout(); |
||||
|
service.shutdown(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testCustomerRestApiLimits() throws Exception { |
||||
|
loginSysAdmin(); |
||||
|
|
||||
|
String customerRestLimit = MESSAGES_LIMIT + ":" + TIME_FOR_LIMIT; |
||||
|
|
||||
|
DefaultTenantProfileConfiguration configurationWithCustomerRestLimits = createTenantProfileConfigurationWithRestLimits(null, customerRestLimit); |
||||
|
|
||||
|
saveTenantProfileWitConfiguration(tenantProfile, configurationWithCustomerRestLimits); |
||||
|
|
||||
|
logout(); |
||||
|
|
||||
|
loginCustomerUser(); |
||||
|
|
||||
|
for (int i = 0; i < MESSAGES_LIMIT; i++) { |
||||
|
doGet("/api/device/types").andExpect(status().isOk()); |
||||
|
} |
||||
|
doGet("/api/device/types").andExpect(status().is4xxClientError()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testTenantRestApiLimits() throws Exception { |
||||
|
loginSysAdmin(); |
||||
|
|
||||
|
String tenantRestLimit = MESSAGES_LIMIT + ":" + TIME_FOR_LIMIT; |
||||
|
|
||||
|
DefaultTenantProfileConfiguration configurationWithTenantRestLimits = createTenantProfileConfigurationWithRestLimits(tenantRestLimit, null); |
||||
|
|
||||
|
saveTenantProfileWitConfiguration(tenantProfile, configurationWithTenantRestLimits); |
||||
|
|
||||
|
logout(); |
||||
|
|
||||
|
loginCustomerUser(); |
||||
|
|
||||
|
for (int i = 0; i < MESSAGES_LIMIT; i++) { |
||||
|
doGet("/api/device/types").andExpect(status().isOk()); |
||||
|
} |
||||
|
doGet("/api/device/types").andExpect(status().is4xxClientError()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testCustomerRestApiLimitsWithAsyncMethod() throws Exception { |
||||
|
loginSysAdmin(); |
||||
|
|
||||
|
String tenantRestLimit = MESSAGES_LIMIT + ":" + TIME_FOR_LIMIT; |
||||
|
|
||||
|
DefaultTenantProfileConfiguration configurationWithTenantRestLimits = createTenantProfileConfigurationWithRestLimits(tenantRestLimit, null); |
||||
|
|
||||
|
saveTenantProfileWitConfiguration(tenantProfile, configurationWithTenantRestLimits); |
||||
|
|
||||
|
logout(); |
||||
|
|
||||
|
loginTenantAdmin(); |
||||
|
|
||||
|
List<ListenableFuture<ResultActions>> attributesRequests = new ArrayList<>(); |
||||
|
|
||||
|
doGet("/api/plugins/telemetry/" + tenantId.getEntityType() + "/" + tenantId.getId().toString() + "/values/attributes").andExpect(status().isOk()); |
||||
|
Thread.sleep(TimeUnit.SECONDS.toMillis(TIME_FOR_LIMIT)); // Wait to initialization for bucket4j
|
||||
|
|
||||
|
for (int i = 0; i < MESSAGES_LIMIT; i++) { |
||||
|
attributesRequests.add(service.submit(() -> doGet("/api/plugins/telemetry/" + tenantId.getEntityType() + "/" + tenantId.getId().toString() + "/values/attributes"))); |
||||
|
} |
||||
|
|
||||
|
List<ResultActions> lists = blockForResponses(attributesRequests); |
||||
|
|
||||
|
for (ResultActions resultActions : lists) { |
||||
|
resultActions.andExpect(status().isOk()); |
||||
|
} |
||||
|
|
||||
|
doGet("/api/plugins/telemetry/" + tenantId.getEntityType() + "/" + tenantId.getId().toString() + "/values/attributes").andExpect(status().is4xxClientError()); |
||||
|
} |
||||
|
|
||||
|
private TenantProfile getDefaultTenantProfile() throws Exception { |
||||
|
|
||||
|
PageLink pageLink = new PageLink(17); |
||||
|
PageData<TenantProfile> pageData = doGetTypedWithPageLink("/api/tenantProfiles?", |
||||
|
new TypeReference<>(){}, pageLink); |
||||
|
Assert.assertFalse(pageData.hasNext()); |
||||
|
Assert.assertEquals(1, pageData.getTotalElements()); |
||||
|
List<TenantProfile> tenantProfiles = new ArrayList<>(pageData.getData()); |
||||
|
|
||||
|
Optional<TenantProfile> optionalDefaultProfile = tenantProfiles.stream().filter(TenantProfile::isDefault).reduce((a, b) -> null); |
||||
|
Assert.assertTrue(optionalDefaultProfile.isPresent()); |
||||
|
|
||||
|
return optionalDefaultProfile.get(); |
||||
|
} |
||||
|
|
||||
|
List<ResultActions> blockForResponses(List<ListenableFuture<ResultActions>> futures) throws ExecutionException { |
||||
|
ListenableFuture<List<ResultActions>> futureOfList = Futures.allAsList(futures); |
||||
|
List<ResultActions> responses; |
||||
|
try { |
||||
|
responses = futureOfList.get(20, TimeUnit.SECONDS); |
||||
|
} catch (TimeoutException | InterruptedException | ExecutionException e) { |
||||
|
responses = new ArrayList<>(); |
||||
|
for (ListenableFuture<ResultActions> future : futures) { |
||||
|
if (future.isDone()) { |
||||
|
responses.add(Uninterruptibles.getUninterruptibly(future)); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return responses; |
||||
|
} |
||||
|
|
||||
|
private DefaultTenantProfileConfiguration createTenantProfileConfigurationWithRestLimits(String tenantLimits, String customerLimits) { |
||||
|
DefaultTenantProfileConfiguration.DefaultTenantProfileConfigurationBuilder builder = DefaultTenantProfileConfiguration.builder(); |
||||
|
builder.tenantServerRestLimitsConfiguration(tenantLimits); |
||||
|
builder.customerServerRestLimitsConfiguration(customerLimits); |
||||
|
return builder.build(); |
||||
|
|
||||
|
} |
||||
|
|
||||
|
private void saveTenantProfileWitConfiguration(TenantProfile tenantProfile, TenantProfileConfiguration tenantProfileConfiguration) { |
||||
|
TenantProfileData tenantProfileData = tenantProfile.getProfileData(); |
||||
|
tenantProfileData.setConfiguration(tenantProfileConfiguration); |
||||
|
TenantProfile savedTenantProfile = doPost("/api/tenantProfile", tenantProfile, TenantProfile.class); |
||||
|
Assert.assertNotNull(savedTenantProfile); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,24 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.system.sql; |
||||
|
|
||||
|
import org.thingsboard.server.dao.service.DaoSqlTest; |
||||
|
import org.thingsboard.server.system.BaseRestApiLimitsTest; |
||||
|
|
||||
|
|
||||
|
@DaoSqlTest |
||||
|
public class RestApiLimitsSqlTest extends BaseRestApiLimitsTest { |
||||
|
} |
||||
@ -0,0 +1,20 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.common.data.script; |
||||
|
|
||||
|
public enum ScriptLanguage { |
||||
|
JS, MVEL |
||||
|
} |
||||
@ -0,0 +1,30 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.common.data.sync.vc; |
||||
|
|
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.Builder; |
||||
|
import lombok.Data; |
||||
|
import lombok.NoArgsConstructor; |
||||
|
|
||||
|
@Data |
||||
|
@AllArgsConstructor |
||||
|
@NoArgsConstructor |
||||
|
@Builder |
||||
|
public class RepositorySettingsInfo { |
||||
|
private boolean configured; |
||||
|
private Boolean readOnly; |
||||
|
} |
||||
@ -0,0 +1,42 @@ |
|||||
|
<!-- |
||||
|
|
||||
|
Copyright © 2016-2022 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. |
||||
|
|
||||
|
--> |
||||
|
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" |
||||
|
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> |
||||
|
<modelVersion>4.0.0</modelVersion> |
||||
|
<parent> |
||||
|
<groupId>org.thingsboard</groupId> |
||||
|
<version>3.4.2-SNAPSHOT</version> |
||||
|
<artifactId>common</artifactId> |
||||
|
</parent> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<artifactId>script</artifactId> |
||||
|
<packaging>pom</packaging> |
||||
|
|
||||
|
<name>Thingsboard Script Invoke Commons</name> |
||||
|
<url>https://thingsboard.io</url> |
||||
|
|
||||
|
<properties> |
||||
|
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> |
||||
|
<main.dir>${basedir}/../..</main.dir> |
||||
|
</properties> |
||||
|
<modules> |
||||
|
<module>script-api</module> |
||||
|
<module>remote-js-client</module> |
||||
|
</modules> |
||||
|
|
||||
|
</project> |
||||
@ -0,0 +1,89 @@ |
|||||
|
<!-- |
||||
|
|
||||
|
Copyright © 2016-2022 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. |
||||
|
|
||||
|
--> |
||||
|
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" |
||||
|
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> |
||||
|
<modelVersion>4.0.0</modelVersion> |
||||
|
<parent> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<version>3.4.2-SNAPSHOT</version> |
||||
|
<artifactId>script</artifactId> |
||||
|
</parent> |
||||
|
<groupId>org.thingsboard.common.script</groupId> |
||||
|
<artifactId>remote-js-client</artifactId> |
||||
|
<packaging>jar</packaging> |
||||
|
|
||||
|
<name>Thingsboard Server JS Client for remote JS execution</name> |
||||
|
<url>https://thingsboard.io</url> |
||||
|
|
||||
|
<properties> |
||||
|
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> |
||||
|
<main.dir>${basedir}/../../..</main.dir> |
||||
|
</properties> |
||||
|
|
||||
|
<dependencies> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<artifactId>queue</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.common.script</groupId> |
||||
|
<artifactId>script-api</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.rule-engine</groupId> |
||||
|
<artifactId>rule-engine-api</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<artifactId>data</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<artifactId>message</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.springframework.boot</groupId> |
||||
|
<artifactId>spring-boot-starter-web</artifactId> |
||||
|
<scope>provided</scope> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.springframework.boot</groupId> |
||||
|
<artifactId>spring-boot-starter-test</artifactId> |
||||
|
<scope>test</scope> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.junit.vintage</groupId> |
||||
|
<artifactId>junit-vintage-engine</artifactId> |
||||
|
<scope>test</scope> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.awaitility</groupId> |
||||
|
<artifactId>awaitility</artifactId> |
||||
|
<scope>test</scope> |
||||
|
</dependency> |
||||
|
</dependencies> |
||||
|
|
||||
|
<distributionManagement> |
||||
|
<repository> |
||||
|
<id>thingsboard-repo-deploy</id> |
||||
|
<name>ThingsBoard Repo Deployment</name> |
||||
|
<url>https://repo.thingsboard.io/artifactory/libs-release-public</url> |
||||
|
</repository> |
||||
|
</distributionManagement> |
||||
|
|
||||
|
</project> |
||||
@ -0,0 +1,121 @@ |
|||||
|
<!-- |
||||
|
|
||||
|
Copyright © 2016-2022 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. |
||||
|
|
||||
|
--> |
||||
|
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" |
||||
|
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> |
||||
|
<modelVersion>4.0.0</modelVersion> |
||||
|
<parent> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<version>3.4.2-SNAPSHOT</version> |
||||
|
<artifactId>script</artifactId> |
||||
|
</parent> |
||||
|
<groupId>org.thingsboard.common.script</groupId> |
||||
|
<artifactId>script-api</artifactId> |
||||
|
<packaging>jar</packaging> |
||||
|
|
||||
|
<name>Thingsboard Server Script invoke API</name> |
||||
|
<url>https://thingsboard.io</url> |
||||
|
|
||||
|
<properties> |
||||
|
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> |
||||
|
<main.dir>${basedir}/../../..</main.dir> |
||||
|
</properties> |
||||
|
|
||||
|
<dependencies> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<artifactId>data</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<artifactId>stats</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard.common</groupId> |
||||
|
<artifactId>util</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.javadelight</groupId> |
||||
|
<artifactId>delight-nashorn-sandbox</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>com.google.code.gson</groupId> |
||||
|
<artifactId>gson</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.slf4j</groupId> |
||||
|
<artifactId>slf4j-api</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.slf4j</groupId> |
||||
|
<artifactId>log4j-over-slf4j</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>ch.qos.logback</groupId> |
||||
|
<artifactId>logback-core</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>ch.qos.logback</groupId> |
||||
|
<artifactId>logback-classic</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.springframework</groupId> |
||||
|
<artifactId>spring-context</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>com.google.guava</groupId> |
||||
|
<artifactId>guava</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.apache.commons</groupId> |
||||
|
<artifactId>commons-lang3</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.thingsboard</groupId> |
||||
|
<artifactId>mvel2</artifactId> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.springframework.boot</groupId> |
||||
|
<artifactId>spring-boot-starter-web</artifactId> |
||||
|
<scope>provided</scope> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.springframework.boot</groupId> |
||||
|
<artifactId>spring-boot-starter-test</artifactId> |
||||
|
<scope>test</scope> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.junit.vintage</groupId> |
||||
|
<artifactId>junit-vintage-engine</artifactId> |
||||
|
<scope>test</scope> |
||||
|
</dependency> |
||||
|
<dependency> |
||||
|
<groupId>org.awaitility</groupId> |
||||
|
<artifactId>awaitility</artifactId> |
||||
|
<scope>test</scope> |
||||
|
</dependency> |
||||
|
</dependencies> |
||||
|
|
||||
|
<distributionManagement> |
||||
|
<repository> |
||||
|
<id>thingsboard-repo-deploy</id> |
||||
|
<name>ThingsBoard Repo Deployment</name> |
||||
|
<url>https://repo.thingsboard.io/artifactory/libs-release-public</url> |
||||
|
</repository> |
||||
|
</distributionManagement> |
||||
|
|
||||
|
</project> |
||||
@ -0,0 +1,288 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api; |
||||
|
|
||||
|
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.extern.slf4j.Slf4j; |
||||
|
import org.springframework.data.util.Pair; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
||||
|
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageReportClient; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageStateClient; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.Executor; |
||||
|
import java.util.concurrent.Executors; |
||||
|
import java.util.concurrent.ScheduledExecutorService; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.concurrent.TimeoutException; |
||||
|
import java.util.concurrent.atomic.AtomicInteger; |
||||
|
|
||||
|
import static java.lang.String.format; |
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class AbstractScriptInvokeService implements ScriptInvokeService { |
||||
|
|
||||
|
protected final Map<UUID, BlockedScriptInfo> disabledScripts = new ConcurrentHashMap<>(); |
||||
|
|
||||
|
private final Optional<TbApiUsageStateClient> apiUsageStateClient; |
||||
|
private final Optional<TbApiUsageReportClient> apiUsageReportClient; |
||||
|
private final AtomicInteger pushedMsgs = new AtomicInteger(0); |
||||
|
private final AtomicInteger invokeMsgs = new AtomicInteger(0); |
||||
|
private final AtomicInteger evalMsgs = new AtomicInteger(0); |
||||
|
protected final AtomicInteger failedMsgs = new AtomicInteger(0); |
||||
|
protected final AtomicInteger timeoutMsgs = new AtomicInteger(0); |
||||
|
|
||||
|
private final FutureCallback<UUID> evalCallback = new ScriptStatCallback<>(evalMsgs, timeoutMsgs, failedMsgs); |
||||
|
private final FutureCallback<Object> invokeCallback = new ScriptStatCallback<>(invokeMsgs, timeoutMsgs, failedMsgs); |
||||
|
|
||||
|
protected ScheduledExecutorService timeoutExecutorService; |
||||
|
|
||||
|
protected AbstractScriptInvokeService(Optional<TbApiUsageStateClient> apiUsageStateClient, Optional<TbApiUsageReportClient> apiUsageReportClient) { |
||||
|
this.apiUsageStateClient = apiUsageStateClient; |
||||
|
this.apiUsageReportClient = apiUsageReportClient; |
||||
|
} |
||||
|
|
||||
|
protected long getMaxEvalRequestsTimeout() { |
||||
|
return getMaxInvokeRequestsTimeout(); |
||||
|
} |
||||
|
|
||||
|
protected abstract long getMaxInvokeRequestsTimeout(); |
||||
|
|
||||
|
protected abstract long getMaxScriptBodySize(); |
||||
|
|
||||
|
protected abstract long getMaxTotalArgsSize(); |
||||
|
|
||||
|
protected abstract long getMaxResultSize(); |
||||
|
|
||||
|
protected abstract int getMaxBlackListDurationSec(); |
||||
|
|
||||
|
protected abstract int getMaxErrors(); |
||||
|
|
||||
|
protected abstract boolean isStatsEnabled(); |
||||
|
|
||||
|
protected abstract String getStatsName(); |
||||
|
|
||||
|
protected abstract Executor getCallbackExecutor(); |
||||
|
|
||||
|
protected abstract boolean isScriptPresent(UUID scriptId); |
||||
|
|
||||
|
protected abstract ListenableFuture<UUID> doEvalScript(TenantId tenantId, ScriptType scriptType, String scriptBody, UUID scriptId, String[] argNames); |
||||
|
|
||||
|
protected abstract TbScriptExecutionTask doInvokeFunction(UUID scriptId, Object[] args); |
||||
|
|
||||
|
protected abstract void doRelease(UUID scriptId) throws Exception; |
||||
|
|
||||
|
public void init() { |
||||
|
if (getMaxEvalRequestsTimeout() > 0 || getMaxInvokeRequestsTimeout() > 0) { |
||||
|
timeoutExecutorService = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("script-timeout")); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void stop() { |
||||
|
if (timeoutExecutorService != null) { |
||||
|
timeoutExecutorService.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void printStats() { |
||||
|
if (isStatsEnabled()) { |
||||
|
int pushed = pushedMsgs.getAndSet(0); |
||||
|
int invoked = invokeMsgs.getAndSet(0); |
||||
|
int evaluated = evalMsgs.getAndSet(0); |
||||
|
int failed = failedMsgs.getAndSet(0); |
||||
|
int timedOut = timeoutMsgs.getAndSet(0); |
||||
|
if (pushed > 0 || invoked > 0 || evaluated > 0 || failed > 0 || timedOut > 0) { |
||||
|
log.info("{}: pushed [{}] received [{}] invoke [{}] eval [{}] failed [{}] timedOut [{}]", |
||||
|
getStatsName(), pushed, invoked + evaluated, invoked, evaluated, failed, timedOut); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<UUID> eval(TenantId tenantId, ScriptType scriptType, String scriptBody, String... argNames) { |
||||
|
if (!apiUsageStateClient.isPresent() || apiUsageStateClient.get().getApiUsageState(tenantId).isJsExecEnabled()) { |
||||
|
if (scriptBodySizeExceeded(scriptBody)) { |
||||
|
return error(format("Script body exceeds maximum allowed size of %s symbols", getMaxScriptBodySize())); |
||||
|
} |
||||
|
UUID scriptId = UUID.randomUUID(); |
||||
|
pushedMsgs.incrementAndGet(); |
||||
|
return withTimeoutAndStatsCallback(scriptId, null, |
||||
|
doEvalScript(tenantId, scriptType, scriptBody, scriptId, argNames), evalCallback, getMaxEvalRequestsTimeout()); |
||||
|
} else { |
||||
|
return error("Script Execution is disabled due to API limits!"); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Object> invokeScript(TenantId tenantId, CustomerId customerId, UUID scriptId, Object... args) { |
||||
|
if (!apiUsageStateClient.isPresent() || apiUsageStateClient.get().getApiUsageState(tenantId).isJsExecEnabled()) { |
||||
|
if (!isScriptPresent(scriptId)) { |
||||
|
return error("No compiled script found for scriptId: [" + scriptId + "]!"); |
||||
|
} |
||||
|
if (!isDisabled(scriptId)) { |
||||
|
if (argsSizeExceeded(args)) { |
||||
|
TbScriptException t = new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, null, new IllegalArgumentException( |
||||
|
format("Script input arguments exceed maximum allowed total args size of %s symbols", getMaxTotalArgsSize()) |
||||
|
)); |
||||
|
return Futures.immediateFailedFuture(handleScriptException(scriptId, null, t)); |
||||
|
} |
||||
|
apiUsageReportClient.ifPresent(client -> client.report(tenantId, customerId, ApiUsageRecordKey.JS_EXEC_COUNT, 1)); |
||||
|
pushedMsgs.incrementAndGet(); |
||||
|
log.trace("InvokeScript uuid {} with timeout {}ms", scriptId, getMaxInvokeRequestsTimeout()); |
||||
|
var task = doInvokeFunction(scriptId, args); |
||||
|
|
||||
|
var resultFuture = Futures.transformAsync(task.getResultFuture(), output -> { |
||||
|
String result = JacksonUtil.toString(output); |
||||
|
if (resultSizeExceeded(result)) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, null, new RuntimeException( |
||||
|
format("Script invocation result exceeds maximum allowed size of %s symbols", getMaxResultSize()) |
||||
|
)); |
||||
|
} |
||||
|
return Futures.immediateFuture(output); |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
|
||||
|
return withTimeoutAndStatsCallback(scriptId, task, resultFuture, invokeCallback, getMaxInvokeRequestsTimeout()); |
||||
|
} else { |
||||
|
String message = "Script invocation is blocked due to maximum error count " |
||||
|
+ getMaxErrors() + ", scriptId " + scriptId + "!"; |
||||
|
log.warn(message); |
||||
|
return error(message); |
||||
|
} |
||||
|
} else { |
||||
|
return error("Script execution is disabled due to API limits!"); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private <T extends V, V> ListenableFuture<T> withTimeoutAndStatsCallback(UUID scriptId, TbScriptExecutionTask task, ListenableFuture<T> future, FutureCallback<V> statsCallback, long timeout) { |
||||
|
if (timeout > 0) { |
||||
|
future = Futures.withTimeout(future, timeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
||||
|
} |
||||
|
Futures.addCallback(future, statsCallback, getCallbackExecutor()); |
||||
|
return Futures.catchingAsync(future, Exception.class, |
||||
|
input -> Futures.immediateFailedFuture(handleScriptException(scriptId, task, input)), |
||||
|
MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
private Throwable handleScriptException(UUID scriptId, TbScriptExecutionTask task, Throwable t) { |
||||
|
boolean timeout = t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException); |
||||
|
if (timeout && task != null) { |
||||
|
task.stop(); |
||||
|
} |
||||
|
boolean blockList = timeout; |
||||
|
String scriptBody = null; |
||||
|
if (t instanceof TbScriptException) { |
||||
|
var scriptException = (TbScriptException) t; |
||||
|
scriptBody = scriptException.getBody(); |
||||
|
var cause = scriptException.getCause(); |
||||
|
switch (scriptException.getErrorCode()) { |
||||
|
case COMPILATION: |
||||
|
log.debug("[{}] Failed to compile script: {}", scriptId, scriptException.getBody(), cause); |
||||
|
break; |
||||
|
case TIMEOUT: |
||||
|
log.debug("[{}] Timeout to execute script: {}", scriptId, scriptException.getBody(), cause); |
||||
|
break; |
||||
|
case OTHER: |
||||
|
case RUNTIME: |
||||
|
log.debug("[{}] Failed to execute script: {}", scriptId, scriptException.getBody(), cause); |
||||
|
break; |
||||
|
} |
||||
|
blockList = timeout || scriptException.getErrorCode() != TbScriptException.ErrorCode.RUNTIME; |
||||
|
} |
||||
|
if (blockList) { |
||||
|
BlockedScriptInfo disableListInfo = disabledScripts.computeIfAbsent(scriptId, key -> new BlockedScriptInfo(getMaxBlackListDurationSec())); |
||||
|
int counter = disableListInfo.incrementAndGet(); |
||||
|
if (log.isDebugEnabled()) { |
||||
|
log.debug("Script has exception counter {} on disabledFunctions for id {}, exception {}, cause {}, scriptBody {}", |
||||
|
counter, scriptId, t, t.getCause(), scriptBody); |
||||
|
} else { |
||||
|
log.warn("Script has exception counter {} on disabledFunctions for id {}, exception {}", |
||||
|
counter, scriptId, t.getMessage()); |
||||
|
} |
||||
|
} |
||||
|
if (timeout) { |
||||
|
return new TimeoutException("Script timeout!"); |
||||
|
} else { |
||||
|
return t; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> release(UUID scriptId) { |
||||
|
if (isScriptPresent(scriptId)) { |
||||
|
try { |
||||
|
disabledScripts.remove(scriptId); |
||||
|
doRelease(scriptId); |
||||
|
} catch (Exception e) { |
||||
|
return Futures.immediateFailedFuture(e); |
||||
|
} |
||||
|
} |
||||
|
return Futures.immediateFuture(null); |
||||
|
} |
||||
|
|
||||
|
private boolean isDisabled(UUID scriptId) { |
||||
|
BlockedScriptInfo errorCount = disabledScripts.get(scriptId); |
||||
|
if (errorCount != null) { |
||||
|
if (errorCount.getExpirationTime() <= System.currentTimeMillis()) { |
||||
|
disabledScripts.remove(scriptId); |
||||
|
return false; |
||||
|
} else { |
||||
|
return errorCount.get() >= getMaxErrors(); |
||||
|
} |
||||
|
} else { |
||||
|
return false; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private boolean scriptBodySizeExceeded(String scriptBody) { |
||||
|
if (getMaxScriptBodySize() <= 0) return false; |
||||
|
return scriptBody.length() > getMaxScriptBodySize(); |
||||
|
} |
||||
|
|
||||
|
private boolean argsSizeExceeded(Object[] args) { |
||||
|
if (getMaxTotalArgsSize() <= 0) return false; |
||||
|
long totalArgsSize = 0; |
||||
|
for (Object arg : args) { |
||||
|
if (arg instanceof CharSequence) { |
||||
|
totalArgsSize += ((CharSequence) arg).length(); |
||||
|
} else { |
||||
|
var str = JacksonUtil.toString(arg); |
||||
|
if (str != null) { |
||||
|
totalArgsSize += str.length(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return totalArgsSize > getMaxTotalArgsSize(); |
||||
|
} |
||||
|
|
||||
|
private boolean resultSizeExceeded(String result) { |
||||
|
if (getMaxResultSize() <= 0) return false; |
||||
|
return result != null && result.length() > getMaxResultSize(); |
||||
|
} |
||||
|
|
||||
|
private <T> ListenableFuture<T> error(String message) { |
||||
|
return Futures.immediateFailedFuture(new RuntimeException(message)); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,44 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api; |
||||
|
|
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.concurrent.atomic.AtomicInteger; |
||||
|
|
||||
|
public class BlockedScriptInfo { |
||||
|
private final long maxScriptBlockDurationMs; |
||||
|
private final AtomicInteger counter; |
||||
|
private long expirationTime; |
||||
|
|
||||
|
BlockedScriptInfo(int maxScriptBlockDuration) { |
||||
|
this.maxScriptBlockDurationMs = TimeUnit.SECONDS.toMillis(maxScriptBlockDuration); |
||||
|
this.counter = new AtomicInteger(0); |
||||
|
} |
||||
|
|
||||
|
public int get() { |
||||
|
return counter.get(); |
||||
|
} |
||||
|
|
||||
|
public int incrementAndGet() { |
||||
|
int result = counter.incrementAndGet(); |
||||
|
expirationTime = System.currentTimeMillis() + maxScriptBlockDurationMs; |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
public long getExpirationTime() { |
||||
|
return expirationTime; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,40 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api; |
||||
|
|
||||
|
import lombok.Getter; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
public class TbScriptException extends RuntimeException { |
||||
|
private static final long serialVersionUID = -1958193538782818284L; |
||||
|
|
||||
|
public static enum ErrorCode {COMPILATION, TIMEOUT, RUNTIME, OTHER} |
||||
|
|
||||
|
@Getter |
||||
|
private final UUID scriptId; |
||||
|
@Getter |
||||
|
private final ErrorCode errorCode; |
||||
|
@Getter |
||||
|
private final String body; |
||||
|
|
||||
|
public TbScriptException(UUID scriptId, ErrorCode errorCode, String body, Exception cause) { |
||||
|
super(cause); |
||||
|
this.scriptId = scriptId; |
||||
|
this.errorCode = errorCode; |
||||
|
this.body = body; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,30 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.Getter; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
|
||||
|
|
||||
|
@RequiredArgsConstructor |
||||
|
public abstract class TbScriptExecutionTask { |
||||
|
|
||||
|
@Getter |
||||
|
private final ListenableFuture<Object> resultFuture; |
||||
|
|
||||
|
public abstract void stop(); |
||||
|
} |
||||
@ -0,0 +1,106 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.js; |
||||
|
|
||||
|
import com.google.common.hash.Hashing; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.Getter; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.data.util.Pair; |
||||
|
import org.thingsboard.script.api.AbstractScriptInvokeService; |
||||
|
import org.thingsboard.script.api.RuleNodeScriptFactory; |
||||
|
import org.thingsboard.script.api.ScriptType; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageReportClient; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageStateClient; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 26.09.18. |
||||
|
*/ |
||||
|
@Slf4j |
||||
|
public abstract class AbstractJsInvokeService extends AbstractScriptInvokeService implements JsInvokeService { |
||||
|
|
||||
|
protected final Map<UUID, JsScriptInfo> scriptInfoMap = new ConcurrentHashMap<>(); |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${js.max_total_args_size:100000}") |
||||
|
private long maxTotalArgsSize; |
||||
|
@Getter |
||||
|
@Value("${js.max_result_size:300000}") |
||||
|
private long maxResultSize; |
||||
|
@Getter |
||||
|
@Value("${js.max_script_body_size:50000}") |
||||
|
private long maxScriptBodySize; |
||||
|
|
||||
|
protected AbstractJsInvokeService(Optional<TbApiUsageStateClient> apiUsageStateClient, Optional<TbApiUsageReportClient> apiUsageReportClient) { |
||||
|
super(apiUsageStateClient, apiUsageReportClient); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected boolean isScriptPresent(UUID scriptId) { |
||||
|
return scriptInfoMap.containsKey(scriptId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected JsScriptExecutionTask doInvokeFunction(UUID scriptId, Object[] args) { |
||||
|
return new JsScriptExecutionTask(doInvokeFunction(scriptId, scriptInfoMap.get(scriptId), args)); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<UUID> doEvalScript(TenantId tenantId, ScriptType scriptType, String scriptBody, UUID scriptId, String[] argNames) { |
||||
|
String scriptHash = hash(tenantId, scriptBody); |
||||
|
String functionName = constructFunctionName(scriptId, scriptHash); |
||||
|
String jsScript = generateJsScript(scriptType, functionName, scriptBody, argNames); |
||||
|
return doEval(scriptId, new JsScriptInfo(scriptHash, functionName), jsScript); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void doRelease(UUID scriptId) throws Exception { |
||||
|
doRelease(scriptId, scriptInfoMap.remove(scriptId)); |
||||
|
} |
||||
|
|
||||
|
protected abstract ListenableFuture<UUID> doEval(UUID scriptId, JsScriptInfo jsInfo, String scriptBody); |
||||
|
|
||||
|
protected abstract ListenableFuture<Object> doInvokeFunction(UUID scriptId, JsScriptInfo jsInfo, Object[] args); |
||||
|
|
||||
|
protected abstract void doRelease(UUID scriptId, JsScriptInfo scriptInfo) throws Exception; |
||||
|
|
||||
|
private String generateJsScript(ScriptType scriptType, String functionName, String scriptBody, String... argNames) { |
||||
|
if (scriptType == ScriptType.RULE_NODE_SCRIPT) { |
||||
|
return RuleNodeScriptFactory.generateRuleNodeScript(functionName, scriptBody, argNames); |
||||
|
} |
||||
|
throw new RuntimeException("No script factory implemented for scriptType: " + scriptType); |
||||
|
} |
||||
|
|
||||
|
protected String constructFunctionName(UUID scriptId, String scriptHash) { |
||||
|
return "invokeInternal_" + scriptId.toString().replace('-', '_'); |
||||
|
} |
||||
|
|
||||
|
protected String hash(TenantId tenantId, String scriptBody) { |
||||
|
return Hashing.murmur3_128().newHasher() |
||||
|
.putLong(tenantId.getId().getMostSignificantBits()) |
||||
|
.putLong(tenantId.getId().getLeastSignificantBits()) |
||||
|
.putUnencodedChars(scriptBody) |
||||
|
.hash().toString(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,28 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.js; |
||||
|
|
||||
|
import org.thingsboard.script.api.ScriptInvokeService; |
||||
|
import org.thingsboard.server.common.data.script.ScriptLanguage; |
||||
|
|
||||
|
public interface JsInvokeService extends ScriptInvokeService { |
||||
|
|
||||
|
@Override |
||||
|
default ScriptLanguage getLanguage() { |
||||
|
return ScriptLanguage.JS; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,31 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.js; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import org.thingsboard.script.api.TbScriptExecutionTask; |
||||
|
|
||||
|
public class JsScriptExecutionTask extends TbScriptExecutionTask { |
||||
|
|
||||
|
public JsScriptExecutionTask(ListenableFuture<Object> resultFuture) { |
||||
|
super(resultFuture); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void stop() { |
||||
|
// do nothing
|
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,179 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.js; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.ListeningExecutorService; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import delight.nashornsandbox.NashornSandbox; |
||||
|
import delight.nashornsandbox.NashornSandboxes; |
||||
|
import lombok.Getter; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
||||
|
import org.springframework.scheduling.annotation.Scheduled; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
||||
|
import org.thingsboard.script.api.TbScriptException; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageReportClient; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageStateClient; |
||||
|
|
||||
|
import javax.annotation.PostConstruct; |
||||
|
import javax.annotation.PreDestroy; |
||||
|
import javax.script.Invocable; |
||||
|
import javax.script.ScriptEngine; |
||||
|
import javax.script.ScriptEngineManager; |
||||
|
import javax.script.ScriptException; |
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.Executor; |
||||
|
import java.util.concurrent.ExecutorService; |
||||
|
import java.util.concurrent.Executors; |
||||
|
import java.util.concurrent.locks.ReentrantLock; |
||||
|
|
||||
|
@Slf4j |
||||
|
@ConditionalOnProperty(prefix = "js", value = "evaluator", havingValue = "local", matchIfMissing = true) |
||||
|
@Service |
||||
|
public class NashornJsInvokeService extends AbstractJsInvokeService { |
||||
|
|
||||
|
private NashornSandbox sandbox; |
||||
|
private ScriptEngine engine; |
||||
|
private ExecutorService monitorExecutorService; |
||||
|
private ListeningExecutorService jsExecutor; |
||||
|
|
||||
|
private final ReentrantLock evalLock = new ReentrantLock(); |
||||
|
|
||||
|
@Value("${js.local.use_js_sandbox}") |
||||
|
private boolean useJsSandbox; |
||||
|
|
||||
|
@Value("${js.local.monitor_thread_pool_size}") |
||||
|
private int monitorThreadPoolSize; |
||||
|
|
||||
|
@Value("${js.local.max_cpu_time}") |
||||
|
private long maxCpuTime; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${js.local.max_errors}") |
||||
|
private int maxErrors; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${js.local.max_black_list_duration_sec:60}") |
||||
|
private int maxBlackListDurationSec; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${js.local.max_requests_timeout:0}") |
||||
|
private long maxInvokeRequestsTimeout; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${js.local.stats.enabled:false}") |
||||
|
private boolean statsEnabled; |
||||
|
|
||||
|
@Value("${js.local.js_thread_pool_size:50}") |
||||
|
private int jsExecutorThreadPoolSize; |
||||
|
|
||||
|
public NashornJsInvokeService(Optional<TbApiUsageStateClient> apiUsageStateClient, Optional<TbApiUsageReportClient> apiUsageReportClient) { |
||||
|
super(apiUsageStateClient, apiUsageReportClient); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected String getStatsName() { |
||||
|
return "Nashorn JS Invoke Stats"; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected Executor getCallbackExecutor() { |
||||
|
return MoreExecutors.directExecutor(); |
||||
|
} |
||||
|
|
||||
|
@Scheduled(fixedDelayString = "${js.local.stats.print_interval_ms:10000}") |
||||
|
public void printStats() { |
||||
|
super.printStats(); |
||||
|
} |
||||
|
|
||||
|
@PostConstruct |
||||
|
public void init() { |
||||
|
super.init(); |
||||
|
jsExecutor = MoreExecutors.listeningDecorator(Executors.newWorkStealingPool(jsExecutorThreadPoolSize)); |
||||
|
if (useJsSandbox) { |
||||
|
sandbox = NashornSandboxes.create(); |
||||
|
monitorExecutorService = ThingsBoardExecutors.newWorkStealingPool(monitorThreadPoolSize, "nashorn-js-monitor"); |
||||
|
sandbox.setExecutor(monitorExecutorService); |
||||
|
sandbox.setMaxCPUTime(maxCpuTime); |
||||
|
sandbox.allowNoBraces(false); |
||||
|
sandbox.allowLoadFunctions(true); |
||||
|
sandbox.setMaxPreparedStatements(30); |
||||
|
} else { |
||||
|
ScriptEngineManager factory = new ScriptEngineManager(); |
||||
|
engine = factory.getEngineByName("nashorn"); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void stop() { |
||||
|
super.stop(); |
||||
|
if (monitorExecutorService != null) { |
||||
|
monitorExecutorService.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<UUID> doEval(UUID scriptId, JsScriptInfo scriptInfo, String jsScript) { |
||||
|
return jsExecutor.submit(() -> { |
||||
|
try { |
||||
|
evalLock.lock(); |
||||
|
try { |
||||
|
if (useJsSandbox) { |
||||
|
sandbox.eval(jsScript); |
||||
|
} else { |
||||
|
engine.eval(jsScript); |
||||
|
} |
||||
|
} finally { |
||||
|
evalLock.unlock(); |
||||
|
} |
||||
|
scriptInfoMap.put(scriptId, scriptInfo); |
||||
|
return scriptId; |
||||
|
} catch (Exception e) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.COMPILATION, jsScript, e); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<Object> doInvokeFunction(UUID scriptId, JsScriptInfo scriptInfo, Object[] args) { |
||||
|
return jsExecutor.submit(() -> { |
||||
|
try { |
||||
|
if (useJsSandbox) { |
||||
|
return sandbox.getSandboxedInvocable().invokeFunction(scriptInfo.getFunctionName(), args); |
||||
|
} else { |
||||
|
return ((Invocable) engine).invokeFunction(scriptInfo.getFunctionName(), args); |
||||
|
} |
||||
|
} catch (ScriptException e) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.RUNTIME, null, e); |
||||
|
} catch (Exception e) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, null, e); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
protected void doRelease(UUID scriptId, JsScriptInfo scriptInfo) throws ScriptException { |
||||
|
if (useJsSandbox) { |
||||
|
sandbox.eval(scriptInfo.getFunctionName() + " = undefined;"); |
||||
|
} else { |
||||
|
engine.eval(scriptInfo.getFunctionName() + " = undefined;"); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,181 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.mvel; |
||||
|
|
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.ListeningExecutorService; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import lombok.Getter; |
||||
|
import lombok.SneakyThrows; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.mvel2.ExecutionContext; |
||||
|
import org.mvel2.MVEL; |
||||
|
import org.mvel2.ParserContext; |
||||
|
import org.mvel2.SandboxedParserConfiguration; |
||||
|
import org.mvel2.SandboxedParserContext; |
||||
|
import org.mvel2.ScriptMemoryOverflowException; |
||||
|
import org.mvel2.optimizers.OptimizerFactory; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
||||
|
import org.springframework.scheduling.annotation.Scheduled; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
||||
|
import org.thingsboard.script.api.AbstractScriptInvokeService; |
||||
|
import org.thingsboard.script.api.ScriptType; |
||||
|
import org.thingsboard.script.api.TbScriptException; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageReportClient; |
||||
|
import org.thingsboard.server.common.stats.TbApiUsageStateClient; |
||||
|
|
||||
|
import javax.annotation.PostConstruct; |
||||
|
import javax.annotation.PreDestroy; |
||||
|
import java.io.Serializable; |
||||
|
import java.util.Map; |
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.Executor; |
||||
|
import java.util.regex.Pattern; |
||||
|
|
||||
|
@Slf4j |
||||
|
@ConditionalOnProperty(prefix = "mvel", value = "enabled", havingValue = "true", matchIfMissing = true) |
||||
|
@Service |
||||
|
public class DefaultMvelInvokeService extends AbstractScriptInvokeService implements MvelInvokeService { |
||||
|
|
||||
|
protected Map<UUID, MvelScript> scriptMap = new ConcurrentHashMap<>(); |
||||
|
private SandboxedParserConfiguration parserConfig; |
||||
|
|
||||
|
private static final Pattern NEW_KEYWORD_PATTERN = Pattern.compile("new\\s"); |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${mvel.max_total_args_size:100000}") |
||||
|
private long maxTotalArgsSize; |
||||
|
@Getter |
||||
|
@Value("${mvel.max_result_size:300000}") |
||||
|
private long maxResultSize; |
||||
|
@Getter |
||||
|
@Value("${mvel.max_script_body_size:50000}") |
||||
|
private long maxScriptBodySize; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${mvel.max_errors:3}") |
||||
|
private int maxErrors; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${mvel.max_black_list_duration_sec:60}") |
||||
|
private int maxBlackListDurationSec; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${mvel.max_requests_timeout:0}") |
||||
|
private long maxInvokeRequestsTimeout; |
||||
|
|
||||
|
@Getter |
||||
|
@Value("${mvel.stats.enabled:false}") |
||||
|
private boolean statsEnabled; |
||||
|
|
||||
|
@Value("${mvel.thread_pool_size:50}") |
||||
|
private int threadPoolSize; |
||||
|
|
||||
|
@Value("${mvel.max_memory_limit_mb:8}") |
||||
|
private long maxMemoryLimitMb; |
||||
|
|
||||
|
private ListeningExecutorService executor; |
||||
|
|
||||
|
protected DefaultMvelInvokeService(Optional<TbApiUsageStateClient> apiUsageStateClient, Optional<TbApiUsageReportClient> apiUsageReportClient) { |
||||
|
super(apiUsageStateClient, apiUsageReportClient); |
||||
|
} |
||||
|
|
||||
|
@Scheduled(fixedDelayString = "${mvel.stats.print_interval_ms:10000}") |
||||
|
public void printStats() { |
||||
|
super.printStats(); |
||||
|
} |
||||
|
|
||||
|
@SneakyThrows |
||||
|
@PostConstruct |
||||
|
public void init() { |
||||
|
super.init(); |
||||
|
OptimizerFactory.setDefaultOptimizer(OptimizerFactory.SAFE_REFLECTIVE); |
||||
|
parserConfig = new SandboxedParserConfiguration(); |
||||
|
parserConfig.addImport("JSON", TbJson.class); |
||||
|
TbUtils.register(parserConfig); |
||||
|
executor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "mvel-executor")); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void destroy() { |
||||
|
if (executor != null) { |
||||
|
executor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected String getStatsName() { |
||||
|
return "MVEL Scripts Stats"; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected Executor getCallbackExecutor() { |
||||
|
return MoreExecutors.directExecutor(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected boolean isScriptPresent(UUID scriptId) { |
||||
|
return scriptMap.containsKey(scriptId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<UUID> doEvalScript(TenantId tenantId, ScriptType scriptType, String scriptBody, UUID scriptId, String[] argNames) { |
||||
|
if (NEW_KEYWORD_PATTERN.matcher(scriptBody).matches()) { |
||||
|
//TODO: output line number and char pos.
|
||||
|
return Futures.immediateFailedFuture(new TbScriptException(scriptId, TbScriptException.ErrorCode.COMPILATION, scriptBody, |
||||
|
new IllegalArgumentException("Keyword 'new' is forbidden!"))); |
||||
|
} |
||||
|
return executor.submit(() -> { |
||||
|
try { |
||||
|
Serializable compiledScript = MVEL.compileExpression(scriptBody, new SandboxedParserContext(parserConfig)); |
||||
|
MvelScript script = new MvelScript(compiledScript, scriptBody, argNames); |
||||
|
scriptMap.put(scriptId, script); |
||||
|
return scriptId; |
||||
|
} catch (Exception e) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.COMPILATION, scriptBody, e); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected MvelScriptExecutionTask doInvokeFunction(UUID scriptId, Object[] args) { |
||||
|
ExecutionContext executionContext = new ExecutionContext(maxMemoryLimitMb * 1024 * 1024); |
||||
|
return new MvelScriptExecutionTask(executionContext, executor.submit(() -> { |
||||
|
MvelScript script = scriptMap.get(scriptId); |
||||
|
if (script == null) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, null, new RuntimeException("Script not found!")); |
||||
|
} |
||||
|
try { |
||||
|
return MVEL.executeTbExpression(script.getCompiledScript(), executionContext, script.createVars(args)); |
||||
|
} catch (ScriptMemoryOverflowException e) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, script.getScriptBody(), new RuntimeException("Script memory overflow!")); |
||||
|
} catch (Exception e) { |
||||
|
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.RUNTIME, script.getScriptBody(), e); |
||||
|
} |
||||
|
})); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void doRelease(UUID scriptId) throws Exception { |
||||
|
scriptMap.remove(scriptId); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,28 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.mvel; |
||||
|
|
||||
|
import org.thingsboard.script.api.ScriptInvokeService; |
||||
|
import org.thingsboard.server.common.data.script.ScriptLanguage; |
||||
|
|
||||
|
public interface MvelInvokeService extends ScriptInvokeService { |
||||
|
|
||||
|
@Override |
||||
|
default ScriptLanguage getLanguage() { |
||||
|
return ScriptLanguage.MVEL; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,41 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.mvel; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
|
||||
|
import java.io.Serializable; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.Map; |
||||
|
|
||||
|
@Data |
||||
|
public class MvelScript { |
||||
|
|
||||
|
private final Serializable compiledScript; |
||||
|
private final String scriptBody; |
||||
|
private final String[] argNames; |
||||
|
|
||||
|
public Map createVars(Object[] args) { |
||||
|
if (args == null || args.length != argNames.length) { |
||||
|
throw new IllegalArgumentException("Invalid number of argument values"); |
||||
|
} |
||||
|
var result = new HashMap<>(); |
||||
|
for (int i = 0; i < argNames.length; i++) { |
||||
|
result.put(argNames[i], args[i]); |
||||
|
} |
||||
|
return result; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,38 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.mvel; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.Data; |
||||
|
import lombok.Getter; |
||||
|
import org.mvel2.ExecutionContext; |
||||
|
import org.thingsboard.script.api.TbScriptExecutionTask; |
||||
|
|
||||
|
|
||||
|
public class MvelScriptExecutionTask extends TbScriptExecutionTask { |
||||
|
|
||||
|
private final ExecutionContext context; |
||||
|
|
||||
|
public MvelScriptExecutionTask(ExecutionContext context, ListenableFuture<Object> resultFuture) { |
||||
|
super(resultFuture); |
||||
|
this.context = context; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void stop(){ |
||||
|
context.stop(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,61 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.mvel; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import org.mvel2.ExecutionContext; |
||||
|
import org.mvel2.util.ArgsRepackUtil; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
|
||||
|
import java.io.IOException; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
|
||||
|
public class TbJson { |
||||
|
|
||||
|
public static String stringify(Object value) { |
||||
|
return value != null ? JacksonUtil.toString(value) : "null"; |
||||
|
} |
||||
|
|
||||
|
public static Object parse(ExecutionContext ctx, String value) throws IOException { |
||||
|
if (value != null) { |
||||
|
JsonNode node = JacksonUtil.toJsonNode(value); |
||||
|
if (node.isObject()) { |
||||
|
return ArgsRepackUtil.repack(ctx, JacksonUtil.convertValue(node, Map.class)); |
||||
|
} else if (node.isArray()) { |
||||
|
return ArgsRepackUtil.repack(ctx, JacksonUtil.convertValue(node, List.class)); |
||||
|
} else if (node.isDouble()) { |
||||
|
return node.doubleValue(); |
||||
|
} else if (node.isLong()) { |
||||
|
return node.longValue(); |
||||
|
} else if (node.isInt()) { |
||||
|
return node.intValue(); |
||||
|
} else if (node.isBoolean()) { |
||||
|
return node.booleanValue(); |
||||
|
} else if (node.isTextual()) { |
||||
|
return node.asText(); |
||||
|
} else if (node.isBinary()) { |
||||
|
return node.binaryValue(); |
||||
|
} else if (node.isNull()) { |
||||
|
return null; |
||||
|
} else { |
||||
|
return node.asText(); |
||||
|
} |
||||
|
} else { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,87 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script.api.mvel; |
||||
|
|
||||
|
import org.mvel2.ExecutionContext; |
||||
|
import org.mvel2.ParserConfiguration; |
||||
|
import org.mvel2.execution.ExecutionArrayList; |
||||
|
import org.mvel2.util.MethodStub; |
||||
|
|
||||
|
import java.io.UnsupportedEncodingException; |
||||
|
import java.util.Base64; |
||||
|
import java.util.List; |
||||
|
|
||||
|
public class TbUtils { |
||||
|
|
||||
|
public static void register(ParserConfiguration parserConfig) throws Exception { |
||||
|
parserConfig.addImport("btoa", new MethodStub(TbUtils.class.getMethod("btoa", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("atob", new MethodStub(TbUtils.class.getMethod("atob", |
||||
|
String.class))); |
||||
|
parserConfig.addImport("bytesToString", new MethodStub(TbUtils.class.getMethod("bytesToString", |
||||
|
List.class))); |
||||
|
parserConfig.addImport("bytesToString", new MethodStub(TbUtils.class.getMethod("bytesToString", |
||||
|
List.class, String.class))); |
||||
|
parserConfig.addImport("stringToBytes", new MethodStub(TbUtils.class.getMethod("stringToBytes", |
||||
|
ExecutionContext.class, String.class))); |
||||
|
parserConfig.addImport("stringToBytes", new MethodStub(TbUtils.class.getMethod("stringToBytes", |
||||
|
ExecutionContext.class, String.class, String.class))); |
||||
|
} |
||||
|
|
||||
|
public static String btoa(String input) { |
||||
|
return new String(Base64.getEncoder().encode(input.getBytes())); |
||||
|
} |
||||
|
|
||||
|
public static String atob(String encoded) { |
||||
|
return new String(Base64.getDecoder().decode(encoded)); |
||||
|
} |
||||
|
|
||||
|
public static String bytesToString(List<Byte> bytesList) { |
||||
|
byte[] bytes = bytesFromList(bytesList); |
||||
|
return new String(bytes); |
||||
|
} |
||||
|
|
||||
|
public static String bytesToString(List<Byte> bytesList, String charsetName) throws UnsupportedEncodingException { |
||||
|
byte[] bytes = bytesFromList(bytesList); |
||||
|
return new String(bytes, charsetName); |
||||
|
} |
||||
|
|
||||
|
public static List<Byte> stringToBytes(ExecutionContext ctx, String str) { |
||||
|
byte[] bytes = str.getBytes(); |
||||
|
return bytesToList(ctx, bytes); |
||||
|
} |
||||
|
|
||||
|
public static List<Byte> stringToBytes(ExecutionContext ctx, String str, String charsetName) throws UnsupportedEncodingException { |
||||
|
byte[] bytes = str.getBytes(charsetName); |
||||
|
return bytesToList(ctx, bytes); |
||||
|
} |
||||
|
|
||||
|
private static byte[] bytesFromList(List<Byte> bytesList) { |
||||
|
byte[] bytes = new byte[bytesList.size()]; |
||||
|
for (int i = 0; i < bytesList.size(); i++) { |
||||
|
bytes[i] = bytesList.get(i); |
||||
|
} |
||||
|
return bytes; |
||||
|
} |
||||
|
|
||||
|
private static List<Byte> bytesToList(ExecutionContext ctx, byte[] bytes) { |
||||
|
List<Byte> list = new ExecutionArrayList<>(ctx); |
||||
|
for (int i = 0; i < bytes.length; i++) { |
||||
|
list.add(bytes[i]); |
||||
|
} |
||||
|
return list; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,25 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.common.stats; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.ApiUsageState; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
|
||||
|
public interface TbApiUsageStateClient { |
||||
|
|
||||
|
ApiUsageState getApiUsageState(TenantId tenantId); |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue