From 4c4923e3f7904e1a2813c9b34a0445ba7998921a Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Thu, 27 Sep 2018 11:59:15 +0300 Subject: [PATCH] Timeout and Args support --- .../service/script/RemoteJsInvokeService.java | 18 +-- application/src/main/proto/jsinvoke.proto | 7 +- .../src/main/resources/thingsboard.yml | 6 +- .../server/kafka/TbJsEvaluator.java | 68 ----------- .../server/kafka/TbRuleEngineEmulator.java | 114 ------------------ 5 files changed, 16 insertions(+), 197 deletions(-) delete mode 100644 common/queue/src/main/java/org/thingsboard/server/kafka/TbJsEvaluator.java delete mode 100644 common/queue/src/main/java/org/thingsboard/server/kafka/TbRuleEngineEmulator.java diff --git a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java index da8d62a231..3507397b7d 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java @@ -47,9 +47,6 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { @Autowired private TbKafkaSettings kafkaSettings; - @Value("${js.remote.use_js_sandbox}") - private boolean useJsSandbox; - @Value("${js.remote.request_topic}") private String requestTopic; @@ -99,8 +96,8 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { } @PreDestroy - public void destroy(){ - if(kafkaTemplate != null){ + public void destroy() { + if (kafkaTemplate != null) { kafkaTemplate.stop(); } } @@ -138,14 +135,19 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { if (scriptBody == null) { return Futures.immediateFailedFuture(new RuntimeException("No script body found for scriptId: [" + scriptId + "]!")); } - JsInvokeProtos.JsInvokeRequest jsRequest = JsInvokeProtos.JsInvokeRequest.newBuilder() + JsInvokeProtos.JsInvokeRequest.Builder jsRequestBuilder = JsInvokeProtos.JsInvokeRequest.newBuilder() .setScriptIdMSB(scriptId.getMostSignificantBits()) .setScriptIdLSB(scriptId.getLeastSignificantBits()) .setFunctionName(functionName) - .setScriptBody(scriptIdToBodysMap.get(scriptId)).build(); + .setTimeout((int) maxRequestsTimeout) + .setScriptBody(scriptIdToBodysMap.get(scriptId)); + + for (int i = 0; i < args.length; i++) { + jsRequestBuilder.setArgs(i, args[i].toString()); + } JsInvokeProtos.RemoteJsRequest jsRequestWrapper = JsInvokeProtos.RemoteJsRequest.newBuilder() - .setInvokeRequest(jsRequest) + .setInvokeRequest(jsRequestBuilder.build()) .build(); ListenableFuture future = kafkaTemplate.post(scriptId.toString(), jsRequestWrapper); diff --git a/application/src/main/proto/jsinvoke.proto b/application/src/main/proto/jsinvoke.proto index 475c7bd82c..d85c3534bb 100644 --- a/application/src/main/proto/jsinvoke.proto +++ b/application/src/main/proto/jsinvoke.proto @@ -19,10 +19,10 @@ package js; option java_package = "org.thingsboard.server.gen.js"; option java_outer_classname = "JsInvokeProtos"; -enum JsInvokeErrorCode{ +enum JsInvokeErrorCode { COMPILATION_ERROR = 0; RUNTIME_ERROR = 1; - CPU_USAGE_ERROR = 2; + TIMEOUT_ERROR = 2; } message RemoteJsRequest { @@ -69,7 +69,8 @@ message JsInvokeRequest { int64 scriptIdLSB = 2; string functionName = 3; string scriptBody = 4; - repeated string args = 5; + int32 timeout = 5; + repeated string args = 6; } message JsInvokeResponse { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index eb40ac506a..37fd6cf687 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -415,7 +415,7 @@ kafka: buffer.memory: "${TB_BUFFER_MEMORY:33554432}" js: - evaluator: "${JS_EVALUATOR:local}" # local/external + evaluator: "${JS_EVALUATOR:local}" # local/remote # Built-in JVM JavaScript environment properties local: # Use Sandboxed (secured) JVM JavaScript environment @@ -428,8 +428,6 @@ js: max_errors: "${LOCAL_JS_SANDBOX_MAX_ERRORS:3}" # Remote JavaScript environment properties remote: - # Use Sandboxed (secured) JVM JavaScript environment - use_js_sandbox: "${USE_REMOTE_JS_SANDBOX:true}" # JS Eval request topic request_topic: "${REMOTE_JS_EVAL_REQUEST_TOPIC:js.eval.requests}" # JS Eval responses topic prefix that is combined with node id @@ -437,7 +435,7 @@ js: # JS Eval max pending requests max_pending_requests: "${REMOTE_JS_MAX_PENDING_REQUESTS:10000}" # JS Eval max request timeout - max_requests_timeout: "${REMOTE_JS_MAX_REQUEST_TIMEOUT:20000}" + max_requests_timeout: "${REMOTE_JS_MAX_REQUEST_TIMEOUT:10000}" # JS response poll interval response_poll_interval: "${REMOTE_JS_RESPONSE_POLL_INTERVAL_MS:25}" # Maximum allowed JavaScript execution errors before JavaScript will be blacklisted diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TbJsEvaluator.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TbJsEvaluator.java deleted file mode 100644 index 92e46b5459..0000000000 --- a/common/queue/src/main/java/org/thingsboard/server/kafka/TbJsEvaluator.java +++ /dev/null @@ -1,68 +0,0 @@ -/** - * Copyright © 2016-2018 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.kafka; - -import lombok.extern.slf4j.Slf4j; -import org.apache.kafka.clients.consumer.ConsumerRecords; -import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.common.header.Header; - -import java.nio.charset.StandardCharsets; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.atomic.LongAdder; - -/** - * Created by ashvayka on 24.09.18. - */ -@Slf4j -public class TbJsEvaluator { - -// public static void main(String[] args) { -// ExecutorService executorService = Executors.newCachedThreadPool(); -// -// TBKafkaConsumerTemplate requestConsumer = new TBKafkaConsumerTemplate(); -// requestConsumer.subscribe("requests"); -// -// LongAdder responseCounter = new LongAdder(); -// TBKafkaProducerTemplate responseProducer = new TBKafkaProducerTemplate(); -// executorService.submit((Runnable) () -> { -// while (true) { -// ConsumerRecords requests = requestConsumer.poll(100); -// requests.forEach(request -> { -// Header header = request.headers().lastHeader("responseTopic"); -// ProducerRecord response = new ProducerRecord<>(new String(header.value(), StandardCharsets.UTF_8), -// request.key(), request.value()); -// responseProducer.send(response); -// responseCounter.add(1); -// }); -// } -// }); -// -// executorService.submit((Runnable) () -> { -// while (true) { -// log.warn("Requests: [{}], Responses: [{}]", responseCounter.longValue(), responseCounter.longValue()); -// try { -// Thread.sleep(1000L); -// } catch (InterruptedException e) { -// e.printStackTrace(); -// } -// } -// }); -// -// } - -} diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TbRuleEngineEmulator.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TbRuleEngineEmulator.java deleted file mode 100644 index ecac2e6796..0000000000 --- a/common/queue/src/main/java/org/thingsboard/server/kafka/TbRuleEngineEmulator.java +++ /dev/null @@ -1,114 +0,0 @@ -/** - * Copyright © 2016-2018 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.kafka; - -import lombok.extern.slf4j.Slf4j; -import org.apache.kafka.clients.admin.CreateTopicsResult; -import org.apache.kafka.clients.admin.NewTopic; -import org.apache.kafka.clients.consumer.ConsumerRecords; -import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.common.header.Header; -import org.apache.kafka.common.header.internals.RecordHeader; - -import java.nio.charset.StandardCharsets; -import java.util.Collections; -import java.util.List; -import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.atomic.LongAdder; - -/** - * Created by ashvayka on 24.09.18. - */ -@Slf4j -public class TbRuleEngineEmulator { -// -// public static void main(String[] args) throws InterruptedException, ExecutionException { -// ConcurrentMap pendingRequestsMap = new ConcurrentHashMap<>(); -// -// ExecutorService executorService = Executors.newCachedThreadPool(); -// -// String responseTopic = "server" + Math.abs((int) (5000.0 * Math.random())); -// try { -// TBKafkaAdmin admin = new TBKafkaAdmin(); -// CreateTopicsResult result = admin.createTopic(new NewTopic(responseTopic, 1, (short) 1)); -// result.all().get(); -// } catch (Exception e) { -// log.warn("Failed to create topic: {}", e.getMessage(), e); -// } -// -// List
headers = Collections.singletonList(new RecordHeader("responseTopic", responseTopic.getBytes(StandardCharsets.UTF_8))); -// -// TBKafkaConsumerTemplate responseConsumer = new TBKafkaConsumerTemplate(); -// TBKafkaProducerTemplate requestProducer = new TBKafkaProducerTemplate(); -// -// LongAdder requestCounter = new LongAdder(); -// LongAdder responseCounter = new LongAdder(); -// -// responseConsumer.subscribe(responseTopic); -// executorService.submit((Runnable) () -> { -// while (true) { -// ConsumerRecords responses = responseConsumer.poll(100); -// responses.forEach(response -> { -// String expectedResponse = pendingRequestsMap.remove(response.key()); -// if (expectedResponse == null) { -// log.error("[{}] Invalid request", response.key()); -// } else if (!expectedResponse.equals(response.value())) { -// log.error("[{}] Invalid response: {} instead of {}", response.key(), response.value(), expectedResponse); -// } -// responseCounter.add(1); -// }); -// } -// }); -// -// executorService.submit((Runnable) () -> { -// int i = 0; -// while (true) { -// String requestId = UUID.randomUUID().toString(); -// String expectedResponse = UUID.randomUUID().toString(); -// pendingRequestsMap.put(requestId, expectedResponse); -// requestProducer.send(new ProducerRecord<>("requests", null, requestId, expectedResponse, headers)); -// requestCounter.add(1); -// i++; -// if (i % 10000 == 0) { -// try { -// Thread.sleep(500L); -// } catch (InterruptedException e) { -// e.printStackTrace(); -// } -// } -// } -// }); -// -// executorService.submit((Runnable) () -> { -// while (true) { -// log.warn("Requests: [{}], Responses: [{}]", requestCounter.longValue(), responseCounter.longValue()); -// try { -// Thread.sleep(1000L); -// } catch (InterruptedException e) { -// e.printStackTrace(); -// } -// } -// }); -// -// Thread.sleep(60000); -// } - -}