Browse Source

Timeouts for Remote JS Executors

pull/2704/head
Andrii Shvaika 7 years ago
parent
commit
4079daabcd
  1. 16
      application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java
  2. 9
      application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java
  3. 16
      application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java
  4. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckAlarmStatusNode.java

16
application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java

@ -18,10 +18,13 @@ package org.thingsboard.server.service.script;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import java.util.Map; import java.util.Map;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicInteger;
/** /**
@ -30,9 +33,22 @@ import java.util.concurrent.atomic.AtomicInteger;
@Slf4j @Slf4j
public abstract class AbstractJsInvokeService implements JsInvokeService { public abstract class AbstractJsInvokeService implements JsInvokeService {
protected ScheduledExecutorService timeoutExecutorService;
protected Map<UUID, String> scriptIdToNameMap = new ConcurrentHashMap<>(); protected Map<UUID, String> scriptIdToNameMap = new ConcurrentHashMap<>();
protected Map<UUID, BlackListInfo> blackListedFunctions = new ConcurrentHashMap<>(); protected Map<UUID, BlackListInfo> blackListedFunctions = new ConcurrentHashMap<>();
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 @Override
public ListenableFuture<UUID> eval(JsScriptType scriptType, String scriptBody, String... argNames) { public ListenableFuture<UUID> eval(JsScriptType scriptType, String scriptBody, String... argNames) {
UUID scriptId = UUID.randomUUID(); UUID scriptId = UUID.randomUUID();

9
application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java

@ -48,7 +48,6 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer
private NashornSandbox sandbox; private NashornSandbox sandbox;
private ScriptEngine engine; private ScriptEngine engine;
private ExecutorService monitorExecutorService; private ExecutorService monitorExecutorService;
private ScheduledExecutorService timeoutExecutorService;
private final AtomicInteger jsPushedMsgs = new AtomicInteger(0); private final AtomicInteger jsPushedMsgs = new AtomicInteger(0);
private final AtomicInteger jsInvokeMsgs = new AtomicInteger(0); private final AtomicInteger jsInvokeMsgs = new AtomicInteger(0);
@ -85,9 +84,7 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer
@PostConstruct @PostConstruct
public void init() { public void init() {
if (maxRequestsTimeout > 0) { super.init(maxRequestsTimeout);
timeoutExecutorService = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("nashorn-js-timeout"));
}
if (useJsSandbox()) { if (useJsSandbox()) {
sandbox = NashornSandboxes.create(); sandbox = NashornSandboxes.create();
monitorExecutorService = Executors.newWorkStealingPool(getMonitorThreadPoolSize()); monitorExecutorService = Executors.newWorkStealingPool(getMonitorThreadPoolSize());
@ -104,12 +101,10 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer
@PreDestroy @PreDestroy
public void stop() { public void stop() {
super.stop();
if (monitorExecutorService != null) { if (monitorExecutorService != null) {
monitorExecutorService.shutdownNow(); monitorExecutorService.shutdownNow();
} }
if (timeoutExecutorService != null) {
timeoutExecutorService.shutdownNow();
}
} }
protected abstract boolean useJsSandbox(); protected abstract boolean useJsSandbox();

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

@ -87,11 +87,13 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
@PostConstruct @PostConstruct
public void init() { public void init() {
super.init(maxRequestsTimeout);
requestTemplate.init(); requestTemplate.init();
} }
@PreDestroy @PreDestroy
public void destroy() { public void destroy() {
super.stop();
if (requestTemplate != null) { if (requestTemplate != null) {
requestTemplate.stop(); requestTemplate.stop();
} }
@ -111,7 +113,9 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
log.trace("Post compile request for scriptId [{}]", scriptId); log.trace("Post compile request for scriptId [{}]", scriptId);
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper));
if (maxRequestsTimeout > 0) {
future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService);
}
kafkaPushedMsgs.incrementAndGet(); kafkaPushedMsgs.incrementAndGet();
Futures.addCallback(future, new FutureCallback<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>>() { Futures.addCallback(future, new FutureCallback<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>>() {
@Override @Override
@ -154,8 +158,8 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
.setTimeout((int) maxRequestsTimeout) .setTimeout((int) maxRequestsTimeout)
.setScriptBody(scriptIdToBodysMap.get(scriptId)); .setScriptBody(scriptIdToBodysMap.get(scriptId));
for (int i = 0; i < args.length; i++) { for (Object arg : args) {
jsRequestBuilder.addArgs(args[i].toString()); jsRequestBuilder.addArgs(arg.toString());
} }
JsInvokeProtos.RemoteJsRequest jsRequestWrapper = JsInvokeProtos.RemoteJsRequest.newBuilder() JsInvokeProtos.RemoteJsRequest jsRequestWrapper = JsInvokeProtos.RemoteJsRequest.newBuilder()
@ -163,6 +167,9 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
.build(); .build();
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper));
if (maxRequestsTimeout > 0) {
future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService);
}
kafkaPushedMsgs.incrementAndGet(); kafkaPushedMsgs.incrementAndGet();
Futures.addCallback(future, new FutureCallback<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>>() { Futures.addCallback(future, new FutureCallback<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>>() {
@Override @Override
@ -203,6 +210,9 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
.build(); .build();
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper));
if (maxRequestsTimeout > 0) {
future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService);
}
JsInvokeProtos.RemoteJsResponse response = future.get().getValue(); JsInvokeProtos.RemoteJsResponse response = future.get().getValue();
JsInvokeProtos.JsReleaseResponse compilationResult = response.getReleaseResponse(); JsInvokeProtos.JsReleaseResponse compilationResult = response.getReleaseResponse();

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckAlarmStatusNode.java

@ -38,7 +38,7 @@ import java.io.IOException;
@Slf4j @Slf4j
@RuleNode( @RuleNode(
type = ComponentType.FILTER, type = ComponentType.FILTER,
name = "checks alarm status", name = "check alarm status",
configClazz = TbCheckAlarmStatusNodeConfig.class, configClazz = TbCheckAlarmStatusNodeConfig.class,
relationTypes = {"True", "False"}, relationTypes = {"True", "False"},
nodeDescription = "Checks alarm status.", nodeDescription = "Checks alarm status.",

Loading…
Cancel
Save