18 changed files with 475 additions and 117 deletions
@ -0,0 +1,46 @@ |
|||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
|
||||
|
version: '2' |
||||
|
|
||||
|
services: |
||||
|
zookeeper: |
||||
|
image: "wurstmeister/zookeeper" |
||||
|
ports: |
||||
|
- "2181" |
||||
|
kafka: |
||||
|
image: "wurstmeister/kafka" |
||||
|
ports: |
||||
|
- "9092:9092" |
||||
|
environment: |
||||
|
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 |
||||
|
KAFKA_LISTENERS: INSIDE://:9093,OUTSIDE://:9092 |
||||
|
KAFKA_ADVERTISED_LISTENERS: INSIDE://:9093,OUTSIDE://${KAFKA_HOSTNAME}:9092 |
||||
|
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INSIDE:PLAINTEXT,OUTSIDE:PLAINTEXT |
||||
|
KAFKA_INTER_BROKER_LISTENER_NAME: INSIDE |
||||
|
KAFKA_CREATE_TOPICS: "${KAFKA_TOPICS}" |
||||
|
KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'false' |
||||
|
depends_on: |
||||
|
- zookeeper |
||||
|
tb-js-executor: |
||||
|
image: "local-maven-build/tb-js-executor:latest" |
||||
|
environment: |
||||
|
TB_KAFKA_SERVERS: kafka:9092 |
||||
|
env_file: |
||||
|
- tb-js-executor.env |
||||
|
depends_on: |
||||
|
- kafka |
||||
@ -0,0 +1,7 @@ |
|||||
|
|
||||
|
REMOTE_JS_EVAL_REQUEST_TOPIC=js.eval.requests |
||||
|
TB_KAFKA_SERVERS=localhost:9092 |
||||
|
LOGGER_LEVEL=debug |
||||
|
LOG_FOLDER=logs |
||||
|
LOGGER_FILENAME=tb-js-executor-%DATE%.log |
||||
|
DOCKER_MODE=true |
||||
@ -0,0 +1,48 @@ |
|||||
|
/* |
||||
|
* 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. |
||||
|
*/ |
||||
|
'use strict'; |
||||
|
|
||||
|
const vm = require('vm'); |
||||
|
|
||||
|
function JsExecutor() { |
||||
|
} |
||||
|
|
||||
|
JsExecutor.prototype.compileScript = function(code) { |
||||
|
return new Promise(function(resolve, reject) { |
||||
|
try { |
||||
|
code = "("+code+")(...args)"; |
||||
|
var script = new vm.Script(code); |
||||
|
resolve(script); |
||||
|
} catch (err) { |
||||
|
reject(err); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
JsExecutor.prototype.executeScript = function(script, args, timeout) { |
||||
|
return new Promise(function(resolve, reject) { |
||||
|
try { |
||||
|
var sandbox = Object.create(null); |
||||
|
sandbox.args = args; |
||||
|
var result = script.runInNewContext(sandbox, {timeout: timeout}); |
||||
|
resolve(result); |
||||
|
} catch (err) { |
||||
|
reject(err); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
module.exports = JsExecutor; |
||||
@ -0,0 +1,226 @@ |
|||||
|
/* |
||||
|
* 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. |
||||
|
*/ |
||||
|
|
||||
|
'use strict'; |
||||
|
|
||||
|
const logger = require('../config/logger')('JsInvokeMessageProcessor'), |
||||
|
Utils = require('./utils'), |
||||
|
js = require('./jsinvoke.proto').js, |
||||
|
KeyedMessage = require('kafka-node').KeyedMessage, |
||||
|
JsExecutor = require('./jsExecutor'); |
||||
|
|
||||
|
function JsInvokeMessageProcessor(producer) { |
||||
|
this.producer = producer; |
||||
|
this.executor = new JsExecutor(); |
||||
|
this.scriptMap = {}; |
||||
|
} |
||||
|
|
||||
|
JsInvokeMessageProcessor.prototype.onJsInvokeMessage = function(message) { |
||||
|
|
||||
|
var requestId; |
||||
|
try { |
||||
|
var request = js.RemoteJsRequest.decode(message.value); |
||||
|
requestId = getRequestId(request); |
||||
|
|
||||
|
logger.debug('[%s] Received request, responseTopic: [%s]', requestId, request.responseTopic); |
||||
|
|
||||
|
if (request.compileRequest) { |
||||
|
this.processCompileRequest(requestId, request.responseTopic, request.compileRequest); |
||||
|
} else if (request.invokeRequest) { |
||||
|
this.processInvokeRequest(requestId, request.responseTopic, request.invokeRequest); |
||||
|
} else if (request.releaseRequest) { |
||||
|
this.processReleaseRequest(requestId, request.responseTopic, request.releaseRequest); |
||||
|
} else { |
||||
|
logger.error('[%s] Unknown request recevied!', requestId); |
||||
|
} |
||||
|
|
||||
|
} catch (err) { |
||||
|
logger.error('[%s] Failed to process request: %s', requestId, err.message); |
||||
|
logger.error(err.stack); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
JsInvokeMessageProcessor.prototype.processCompileRequest = function(requestId, responseTopic, compileRequest) { |
||||
|
var scriptId = getScriptId(compileRequest); |
||||
|
logger.debug('[%s] Processing compile request, scriptId: [%s]', requestId, scriptId); |
||||
|
|
||||
|
this.executor.compileScript(compileRequest.scriptBody).then( |
||||
|
(script) => { |
||||
|
this.scriptMap[scriptId] = script; |
||||
|
var compileResponse = createCompileResponse(scriptId, true); |
||||
|
logger.debug('[%s] Sending success compile response, scriptId: [%s]', requestId, scriptId); |
||||
|
this.sendResponse(requestId, responseTopic, scriptId, compileResponse); |
||||
|
}, |
||||
|
(err) => { |
||||
|
var compileResponse = createCompileResponse(scriptId, false, js.JsInvokeErrorCode.COMPILATION_ERROR, err); |
||||
|
logger.debug('[%s] Sending failed compile response, scriptId: [%s]', requestId, scriptId); |
||||
|
this.sendResponse(requestId, responseTopic, scriptId, compileResponse); |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
JsInvokeMessageProcessor.prototype.processInvokeRequest = function(requestId, responseTopic, invokeRequest) { |
||||
|
var scriptId = getScriptId(invokeRequest); |
||||
|
logger.debug('[%s] Processing invoke request, scriptId: [%s]', requestId, scriptId); |
||||
|
this.getOrCompileScript(scriptId, invokeRequest.scriptBody).then( |
||||
|
(script) => { |
||||
|
this.executor.executeScript(script, invokeRequest.args, invokeRequest.timeout).then( |
||||
|
(result) => { |
||||
|
var invokeResponse = createInvokeResponse(result, true); |
||||
|
logger.debug('[%s] Sending success invoke response, scriptId: [%s]', requestId, scriptId); |
||||
|
this.sendResponse(requestId, responseTopic, scriptId, null, invokeResponse); |
||||
|
}, |
||||
|
(err) => { |
||||
|
var errorCode; |
||||
|
if (err.message.includes('Script execution timed out')) { |
||||
|
errorCode = js.JsInvokeErrorCode.TIMEOUT_ERROR; |
||||
|
} else { |
||||
|
errorCode = js.JsInvokeErrorCode.RUNTIME_ERROR; |
||||
|
} |
||||
|
var invokeResponse = createInvokeResponse("", false, errorCode, err); |
||||
|
logger.debug('[%s] Sending failed invoke response, scriptId: [%s], errorCode: [%s]', requestId, scriptId, errorCode); |
||||
|
this.sendResponse(requestId, responseTopic, scriptId, null, invokeResponse); |
||||
|
} |
||||
|
) |
||||
|
}, |
||||
|
(err) => { |
||||
|
var invokeResponse = createInvokeResponse("", false, js.JsInvokeErrorCode.COMPILATION_ERROR, err); |
||||
|
logger.debug('[%s] Sending failed invoke response, scriptId: [%s], errorCode: [%s]', requestId, scriptId, js.JsInvokeErrorCode.COMPILATION_ERROR); |
||||
|
this.sendResponse(requestId, responseTopic, scriptId, null, invokeResponse); |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
JsInvokeMessageProcessor.prototype.processReleaseRequest = function(requestId, responseTopic, releaseRequest) { |
||||
|
var scriptId = getScriptId(releaseRequest); |
||||
|
logger.debug('[%s] Processing release request, scriptId: [%s]', requestId, scriptId); |
||||
|
if (this.scriptMap[scriptId]) { |
||||
|
delete this.scriptMap[scriptId]; |
||||
|
} |
||||
|
var releaseResponse = createReleaseResponse(scriptId, true); |
||||
|
logger.debug('[%s] Sending success release response, scriptId: [%s]', requestId, scriptId); |
||||
|
this.sendResponse(requestId, responseTopic, scriptId, null, null, releaseResponse); |
||||
|
} |
||||
|
|
||||
|
JsInvokeMessageProcessor.prototype.sendResponse = function (requestId, responseTopic, scriptId, compileResponse, invokeResponse, releaseResponse) { |
||||
|
var remoteResponse = createRemoteResponse(requestId, compileResponse, invokeResponse, releaseResponse); |
||||
|
var rawResponse = js.RemoteJsResponse.encode(remoteResponse).finish(); |
||||
|
const message = new KeyedMessage(scriptId, rawResponse); |
||||
|
const payloads = [ { topic: responseTopic, messages: message, key: scriptId } ]; |
||||
|
this.producer.send(payloads, function (err, data) { |
||||
|
if (err) { |
||||
|
logger.error('[%s] Failed to send response to kafka: %s', requestId, err.message); |
||||
|
logger.error(err.stack); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scriptBody) { |
||||
|
var self = this; |
||||
|
return new Promise(function(resolve, reject) { |
||||
|
if (self.scriptMap[scriptId]) { |
||||
|
resolve(self.scriptMap[scriptId]); |
||||
|
} else { |
||||
|
self.executor.compileScript(scriptBody).then( |
||||
|
(script) => { |
||||
|
self.scriptMap[scriptId] = script; |
||||
|
resolve(script); |
||||
|
}, |
||||
|
(err) => { |
||||
|
reject(err); |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
function createRemoteResponse(requestId, compileResponse, invokeResponse, releaseResponse) { |
||||
|
const requestIdBits = Utils.UUIDToBits(requestId); |
||||
|
return js.RemoteJsResponse.create( |
||||
|
{ |
||||
|
requestIdMSB: requestIdBits[0], |
||||
|
requestIdLSB: requestIdBits[1], |
||||
|
compileResponse: compileResponse, |
||||
|
invokeResponse: invokeResponse, |
||||
|
releaseResponse: releaseResponse |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
function createCompileResponse(scriptId, success, errorCode, err) { |
||||
|
const scriptIdBits = Utils.UUIDToBits(scriptId); |
||||
|
return js.JsCompileResponse.create( |
||||
|
{ |
||||
|
errorCode: errorCode, |
||||
|
success: success, |
||||
|
errorDetails: parseJsErrorDetails(err), |
||||
|
scriptIdMSB: scriptIdBits[0], |
||||
|
scriptIdLSB: scriptIdBits[1] |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
function createInvokeResponse(result, success, errorCode, err) { |
||||
|
return js.JsInvokeResponse.create( |
||||
|
{ |
||||
|
errorCode: errorCode, |
||||
|
success: success, |
||||
|
errorDetails: parseJsErrorDetails(err), |
||||
|
result: result |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
function createReleaseResponse(scriptId, success) { |
||||
|
const scriptIdBits = Utils.UUIDToBits(scriptId); |
||||
|
return js.JsReleaseResponse.create( |
||||
|
{ |
||||
|
success: success, |
||||
|
scriptIdMSB: scriptIdBits[0], |
||||
|
scriptIdLSB: scriptIdBits[1] |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
function parseJsErrorDetails(err) { |
||||
|
if (!err) { |
||||
|
return ''; |
||||
|
} |
||||
|
var details = err.name + ': ' + err.message; |
||||
|
if (err.stack) { |
||||
|
var lines = err.stack.split('\n'); |
||||
|
if (lines && lines.length) { |
||||
|
var line = lines[0]; |
||||
|
var splitted = line.split(':'); |
||||
|
if (splitted && splitted.length === 2) { |
||||
|
if (!isNaN(splitted[1])) { |
||||
|
details += ' in at line number ' + splitted[1]; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return details; |
||||
|
} |
||||
|
|
||||
|
function getScriptId(request) { |
||||
|
return Utils.toUUIDString(request.scriptIdMSB, request.scriptIdLSB); |
||||
|
} |
||||
|
|
||||
|
function getRequestId(request) { |
||||
|
return Utils.toUUIDString(request.requestIdMSB, request.requestIdLSB); |
||||
|
} |
||||
|
|
||||
|
module.exports = JsInvokeMessageProcessor; |
||||
@ -1,78 +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. |
|
||||
*/ |
|
||||
|
|
||||
'use strict'; |
|
||||
|
|
||||
const logger = require('../config/logger')('JsMessageConsumer'); |
|
||||
const Utils = require('./Utils'); |
|
||||
const js = require('./jsinvoke.proto').js; |
|
||||
const KeyedMessage = require('kafka-node').KeyedMessage; |
|
||||
|
|
||||
|
|
||||
exports.onJsInvokeMessage = function(message, producer) { |
|
||||
|
|
||||
logger.info('Received message: %s', JSON.stringify(message)); |
|
||||
|
|
||||
var request = js.RemoteJsRequest.decode(message.value); |
|
||||
|
|
||||
logger.info('Received request: %s', JSON.stringify(request)); |
|
||||
|
|
||||
var requestId = getRequestId(request); |
|
||||
|
|
||||
logger.info('Received request, responseTopic: [%s]; requestId: [%s]', request.responseTopic, requestId); |
|
||||
|
|
||||
if (request.compileRequest) { |
|
||||
var scriptId = getScriptId(request.compileRequest); |
|
||||
|
|
||||
logger.info('Received compile request, scriptId: [%s]', scriptId); |
|
||||
|
|
||||
var compileResponse = js.JsCompileResponse.create( |
|
||||
{ |
|
||||
errorCode: js.JsInvokeErrorCode.COMPILATION_ERROR, |
|
||||
success: false, |
|
||||
errorDetails: 'Not Implemented!', |
|
||||
scriptIdLSB: request.compileRequest.scriptIdLSB, |
|
||||
scriptIdMSB: request.compileRequest.scriptIdMSB |
|
||||
} |
|
||||
); |
|
||||
const requestIdBits = Utils.UUIDToBits(requestId); |
|
||||
var response = js.RemoteJsResponse.create( |
|
||||
{ |
|
||||
requestIdMSB: requestIdBits[0], |
|
||||
requestIdLSB: requestIdBits[1], |
|
||||
compileResponse: compileResponse |
|
||||
} |
|
||||
); |
|
||||
var rawResponse = js.RemoteJsResponse.encode(response).finish(); |
|
||||
sendMessage(producer, rawResponse, request.responseTopic, scriptId); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
function sendMessage(producer, rawMessage, responseTopic, scriptId) { |
|
||||
const message = new KeyedMessage(scriptId, rawMessage); |
|
||||
const payloads = [ { topic: responseTopic, messages: rawMessage, key: scriptId } ]; |
|
||||
producer.send(payloads, function (err, data) { |
|
||||
console.log(data); |
|
||||
}); |
|
||||
} |
|
||||
|
|
||||
function getScriptId(request) { |
|
||||
return Utils.toUUIDString(request.scriptIdMSB, request.scriptIdLSB); |
|
||||
} |
|
||||
|
|
||||
function getRequestId(request) { |
|
||||
return Utils.toUUIDString(request.requestIdMSB, request.requestIdLSB); |
|
||||
} |
|
||||
@ -0,0 +1,26 @@ |
|||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
FROM debian:stretch |
||||
|
|
||||
|
COPY start-js-executor.sh ${pkg.name}.deb /tmp/ |
||||
|
|
||||
|
RUN chmod a+x /tmp/*.sh \ |
||||
|
&& mv /tmp/start-js-executor.sh /usr/bin |
||||
|
|
||||
|
RUN dpkg -i /tmp/${pkg.name}.deb |
||||
|
|
||||
|
CMD ["start-js-executor.sh"] |
||||
@ -0,0 +1,29 @@ |
|||||
|
#!/bin/bash |
||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
|
||||
|
echo "Starting '${project.name}' ..." |
||||
|
|
||||
|
CONF_FOLDER="${pkg.installFolder}/conf" |
||||
|
|
||||
|
mainfile=${pkg.installFolder}/bin/${pkg.name} |
||||
|
configfile=${pkg.name}.conf |
||||
|
identity=${pkg.name} |
||||
|
|
||||
|
source "${CONF_FOLDER}/${configfile}" |
||||
|
|
||||
|
su -s /bin/sh -c "$mainfile" |
||||
Loading…
Reference in new issue