From 0c886656543987f27d59c2c053c11677de147a8b Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 25 May 2021 19:44:34 +0300 Subject: [PATCH 01/21] js-executor: scriptMap refactored from Object to the Map() --- .../api/jsInvokeMessageProcessor.js | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index 2964291a30..a760e21e47 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -34,10 +34,9 @@ const slowQueryLogBody = config.get('script.slow_query_log_body') === 'true'; const {performance} = require('perf_hooks'); function JsInvokeMessageProcessor(producer) { - console.log("Producer:", producer); this.producer = producer; this.executor = new JsExecutor(useSandbox); - this.scriptMap = {}; + this.scriptMap = new Map(); this.scriptIds = []; this.executedScriptsCounter = 0; } @@ -157,12 +156,12 @@ JsInvokeMessageProcessor.prototype.processInvokeRequest = function(requestId, re JsInvokeMessageProcessor.prototype.processReleaseRequest = function(requestId, responseTopic, headers, releaseRequest) { var scriptId = getScriptId(releaseRequest); logger.debug('[%s] Processing release request, scriptId: [%s]', requestId, scriptId); - if (this.scriptMap[scriptId]) { + if (this.scriptMap.has(scriptId)) { var index = this.scriptIds.indexOf(scriptId); if (index > -1) { this.scriptIds.splice(index, 1); } - delete this.scriptMap[scriptId]; + this.scriptMap.delete(scriptId); } var releaseResponse = createReleaseResponse(scriptId, true); logger.debug('[%s] Sending success release response, scriptId: [%s]', requestId, scriptId); @@ -189,8 +188,8 @@ JsInvokeMessageProcessor.prototype.sendResponse = function (requestId, responseT JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scriptBody) { var self = this; return new Promise(function(resolve, reject) { - if (self.scriptMap[scriptId]) { - resolve(self.scriptMap[scriptId]); + if (self.scriptMap.has(scriptId)) { + resolve(self.scriptMap.get(scriptId)); } else { self.executor.compileScript(scriptBody).then( (script) => { @@ -206,16 +205,17 @@ JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scrip } JsInvokeMessageProcessor.prototype.cacheScript = function(scriptId, script) { - if (!this.scriptMap[scriptId]) { + if (!this.scriptMap.has(scriptId)) { this.scriptIds.push(scriptId); while (this.scriptIds.length > maxActiveScripts) { logger.info('Active scripts count [%s] exceeds maximum limit [%s]', this.scriptIds.length, maxActiveScripts); const prevScriptId = this.scriptIds.shift(); logger.info('Removing active script with id [%s]', prevScriptId); - delete this.scriptMap[prevScriptId]; + this.scriptMap.delete(prevScriptId); } } - this.scriptMap[scriptId] = script; + this.scriptMap.set(scriptId, script); + logger.info("scriptMap size is [%s]", this.scriptMap.size); } function createRemoteResponse(requestId, compileResponse, invokeResponse, releaseResponse) { From d729d9ee959963b570aeeaf27787cdd6b6f44bdf Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 27 May 2021 14:50:42 +0300 Subject: [PATCH 02/21] js-executor: instrumentation event for producer and consumer to define the exact flow how to Kafka works without batches (for debug only) --- msa/js-executor/queue/kafkaTemplate.js | 24 ++++++++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 2f4fd8751d..1fd21185d0 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -114,6 +114,30 @@ function KafkaProducer() { consumer = kafkaClient.consumer({groupId: 'js-executor-group'}); producer = kafkaClient.producer(); + + const { CONNECT } = producer.events; + const removeListenerC = producer.on(CONNECT, e => logger.info(`producer CONNECT`)); + const { DISCONNECT } = producer.events; + const removeListenerD = producer.on(DISCONNECT, e => logger.info(`producer DISCONNECT`)); + const { REQUEST } = producer.events; + const removeListenerR = producer.on(REQUEST, e => logger.info(`producer REQUEST ${e.payload.broker}`)); + const { REQUEST_TIMEOUT } = producer.events; + const removeListenerRT = producer.on(REQUEST_TIMEOUT, e => logger.info(`producer REQUEST_TIMEOUT ${e.payload.broker}`)); + const { REQUEST_QUEUE_SIZE } = producer.events; + const removeListenerRQS = producer.on(REQUEST_QUEUE_SIZE, e => logger.info(`producer REQUEST_QUEUE_SIZE ${e.payload.broker} size ${e.queueSize}`)); + + const removeListeners = {} + const { FETCH_START } = consumer.events; + removeListeners[FETCH_START] = consumer.on(FETCH_START, e => logger.info(`consumer FETCH_START`)); + const { FETCH } = consumer.events; + removeListeners[FETCH] = consumer.on(FETCH, e => logger.info(`consumer FETCH numberOfBatches ${e.payload.numberOfBatches} duration ${e.payload.duration}`)); + const { START_BATCH_PROCESS } = consumer.events; + removeListeners[START_BATCH_PROCESS] = consumer.on(START_BATCH_PROCESS, e => logger.info(`consumer START_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); + const { END_BATCH_PROCESS } = consumer.events; + removeListeners[END_BATCH_PROCESS] = consumer.on(END_BATCH_PROCESS, e => logger.info(`consumer END_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); + const { COMMIT_OFFSETS } = consumer.events; + removeListeners[COMMIT_OFFSETS] = consumer.on(COMMIT_OFFSETS, e => logger.info(`consumer COMMIT_OFFSETS topics ${e.payload.topics}`)); + const messageProcessor = new JsInvokeMessageProcessor(new KafkaProducer()); await consumer.connect(); await producer.connect(); From 35e2ff99c3d3f7a488fdc4d3fb3375daa06d41f9 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 27 May 2021 14:55:10 +0300 Subject: [PATCH 03/21] js-executor: send messages as batch --- .../api/jsInvokeMessageProcessor.js | 25 ++++---- msa/js-executor/queue/kafkaTemplate.js | 63 ++++++++++++++----- 2 files changed, 63 insertions(+), 25 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index a760e21e47..867dd82b39 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -172,17 +172,20 @@ JsInvokeMessageProcessor.prototype.sendResponse = function (requestId, responseT var tStartSending = performance.now(); var remoteResponse = createRemoteResponse(requestId, compileResponse, invokeResponse, releaseResponse); var rawResponse = Buffer.from(JSON.stringify(remoteResponse), 'utf8'); - this.producer.send(responseTopic, scriptId, rawResponse, headers).then( - () => { - logger.debug('[%s] Response sent to queue, took [%s]ms, scriptId: [%s]', requestId, (performance.now() - tStartSending), scriptId); - }, - (err) => { - if (err) { - logger.error('[%s] Failed to send response to queue: %s', requestId, err.message); - logger.error(err.stack); - } - } - ); + logger.debug('[%s] Sending response to queue, scriptId: [%s]', requestId, scriptId); + this.producer.send(responseTopic, scriptId, rawResponse, headers); +//TODO put error msg for other queues implementation except Kafka +// .then( +// () => { +// logger.info('[%s] Response sent to queue, took [%s]ms, scriptId: [%s]', requestId, (performance.now() - tStartSending), scriptId); +// }, +// (err) => { +// if (err) { +// logger.error('[%s] Failed to send response to queue: %s', requestId, err.message); +// logger.error(err.stack); +// } +// } +// ); } JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scriptBody) { diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 1fd21185d0..da46a45301 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -33,24 +33,54 @@ let producer; const configEntries = []; +let topicMessages = []; +let loopSend; + function KafkaProducer() { - this.send = async (responseTopic, scriptId, rawResponse, headers) => { - return producer.send( - { - topic: responseTopic, - acks: acks, - compression: compressionType, - messages: [ - { - key: scriptId, - value: rawResponse, - headers: headers.data - } - ] - }); + this.send = (responseTopic, scriptId, rawResponse, headers) => { + logger.debug('Pending queue response, scriptId: [%s]', scriptId); + const message = { + topic: responseTopic, + messages: [{ + key: scriptId, + value: rawResponse, + headers: headers.data + }] + }; + + topicMessages.push(message); + return {}; } } +function sendLoopFunction() { + loopSend = setInterval(sendProducerMsg, 200); +} + +function sendProducerMsg() { + if (topicMessages.length > 0) { + logger.info('sendProducerMsg from queue response, lenght: [%s]', topicMessages.length ); + const messagesToSend = topicMessages; + topicMessages = []; + producer.sendBatch({ + topicMessages: messagesToSend, + acks: acks, + compression: compressionType + }).then( + () => { + logger.info('Response sent to kafka, length: [%s]', messagesToSend.length ); + }, + (err) => { + if (err) { + logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length ); + topicMessages = messagesToSend.concat(topicMessages); + logger.error(err.stack); + } + } + ); + } +} + (async () => { try { logger.info('Starting ThingsBoard JavaScript Executor Microservice...'); @@ -141,6 +171,7 @@ function KafkaProducer() { const messageProcessor = new JsInvokeMessageProcessor(new KafkaProducer()); await consumer.connect(); await producer.connect(); + sendLoopFunction(); await consumer.subscribe({topic: requestTopic}); logger.info('Started ThingsBoard JavaScript Executor Microservice.'); @@ -221,6 +252,10 @@ async function disconnectProducer() { var _producer = producer; producer = null; try { + logger.info('Stopping loop...'); + //TODO: send handle msg + clearInterval(loopSend); + sendProducerMsg(); await _producer.disconnect(); logger.info('Kafka Producer stopped.'); } catch (e) { From 0970ce65b448f622056c070831aa1625f1e81f36 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 27 May 2021 18:45:29 +0300 Subject: [PATCH 04/21] js-executor: maxBatchSize --- msa/js-executor/queue/kafkaTemplate.js | 88 ++++++++++++++++---------- 1 file changed, 53 insertions(+), 35 deletions(-) diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index da46a45301..a061215b28 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -26,6 +26,9 @@ const acks = Number(config.get('kafka.acks')); const requestTimeout = Number(config.get('kafka.requestTimeout')); const compressionType = (config.get('kafka.requestTimeout') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; +const linger = 5; //milliseconds //TODO move to the config +const maxBatchSize = 10; //max messages in batch //TODO move to the config + let kafkaClient; let kafkaAdmin; let consumer; @@ -33,8 +36,8 @@ let producer; const configEntries = []; -let topicMessages = []; -let loopSend; +let batchMessages = []; +let sendLoopInstance; function KafkaProducer() { this.send = (responseTopic, scriptId, rawResponse, headers) => { @@ -48,36 +51,51 @@ function KafkaProducer() { }] }; - topicMessages.push(message); + pushMessageToSendLater(message); return {}; } } -function sendLoopFunction() { - loopSend = setInterval(sendProducerMsg, 200); +function pushMessageToSendLater(message) { + batchMessages.push(message); + if (batchMessages.length >= maxBatchSize) { + sendMessagesAsBatch(); + sendLoopWithLinger(); //reset loop function and reschedule new linger + } +} + +function sendLoopWithLinger() { + if (sendLoopInstance) { + logger.debug("Clear sendLoop scheduler. Starting new send loop with linger [%s]", linger); + clearInterval(sendLoopInstance); + } else { + logger.debug("Starting new send loop with linger [%s]", linger) + } + sendLoopInstance = setInterval(sendMessagesAsBatch, linger); } -function sendProducerMsg() { - if (topicMessages.length > 0) { - logger.info('sendProducerMsg from queue response, lenght: [%s]', topicMessages.length ); - const messagesToSend = topicMessages; - topicMessages = []; +function sendMessagesAsBatch() { + if (batchMessages.length > 0) { + logger.info('sendMessagesAsBatch, lenght: [%s]', batchMessages.length ); + const messagesToSend = batchMessages; + batchMessages = []; producer.sendBatch({ topicMessages: messagesToSend, acks: acks, compression: compressionType }).then( - () => { - logger.info('Response sent to kafka, length: [%s]', messagesToSend.length ); - }, - (err) => { - if (err) { - logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length ); - topicMessages = messagesToSend.concat(topicMessages); - logger.error(err.stack); - } - } - ); + () => { + logger.info('Response sent to kafka, length: [%s]', messagesToSend.length ); + }, + (err) => { + logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length ); + batchMessages = messagesToSend.concat(batchMessages); + logger.error(err.stack); + } + ); + + } else { + //logger.debug("nothing to send"); } } @@ -156,22 +174,22 @@ function sendProducerMsg() { const { REQUEST_QUEUE_SIZE } = producer.events; const removeListenerRQS = producer.on(REQUEST_QUEUE_SIZE, e => logger.info(`producer REQUEST_QUEUE_SIZE ${e.payload.broker} size ${e.queueSize}`)); - const removeListeners = {} - const { FETCH_START } = consumer.events; - removeListeners[FETCH_START] = consumer.on(FETCH_START, e => logger.info(`consumer FETCH_START`)); - const { FETCH } = consumer.events; - removeListeners[FETCH] = consumer.on(FETCH, e => logger.info(`consumer FETCH numberOfBatches ${e.payload.numberOfBatches} duration ${e.payload.duration}`)); - const { START_BATCH_PROCESS } = consumer.events; - removeListeners[START_BATCH_PROCESS] = consumer.on(START_BATCH_PROCESS, e => logger.info(`consumer START_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); - const { END_BATCH_PROCESS } = consumer.events; - removeListeners[END_BATCH_PROCESS] = consumer.on(END_BATCH_PROCESS, e => logger.info(`consumer END_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); - const { COMMIT_OFFSETS } = consumer.events; - removeListeners[COMMIT_OFFSETS] = consumer.on(COMMIT_OFFSETS, e => logger.info(`consumer COMMIT_OFFSETS topics ${e.payload.topics}`)); +// const removeListeners = {} +// const { FETCH_START } = consumer.events; +// removeListeners[FETCH_START] = consumer.on(FETCH_START, e => logger.info(`consumer FETCH_START`)); +// const { FETCH } = consumer.events; +// removeListeners[FETCH] = consumer.on(FETCH, e => logger.info(`consumer FETCH numberOfBatches ${e.payload.numberOfBatches} duration ${e.payload.duration}`)); +// const { START_BATCH_PROCESS } = consumer.events; +// removeListeners[START_BATCH_PROCESS] = consumer.on(START_BATCH_PROCESS, e => logger.info(`consumer START_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); +// const { END_BATCH_PROCESS } = consumer.events; +// removeListeners[END_BATCH_PROCESS] = consumer.on(END_BATCH_PROCESS, e => logger.info(`consumer END_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); +// const { COMMIT_OFFSETS } = consumer.events; +// removeListeners[COMMIT_OFFSETS] = consumer.on(COMMIT_OFFSETS, e => logger.info(`consumer COMMIT_OFFSETS topics ${e.payload.topics}`)); const messageProcessor = new JsInvokeMessageProcessor(new KafkaProducer()); await consumer.connect(); await producer.connect(); - sendLoopFunction(); + sendLoopWithLinger(); await consumer.subscribe({topic: requestTopic}); logger.info('Started ThingsBoard JavaScript Executor Microservice.'); @@ -254,8 +272,8 @@ async function disconnectProducer() { try { logger.info('Stopping loop...'); //TODO: send handle msg - clearInterval(loopSend); - sendProducerMsg(); + clearInterval(sendLoopInstance); + sendMessagesAsBatch(); await _producer.disconnect(); logger.info('Kafka Producer stopped.'); } catch (e) { From 93bea7020584343511b1a561bb861fd010df037e Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 28 May 2021 09:41:35 +0300 Subject: [PATCH 05/21] js-executor: added parameters for producer TB_KAFKA_BATCH_SIZE and TB_KAFKA_LINGER_MS; added print stats frequency SCRIPT_STAT_PRINT_FREQUENCY --- .../api/jsInvokeMessageProcessor.js | 17 +++++--- .../config/custom-environment-variables.yml | 3 ++ msa/js-executor/config/default.yml | 7 ++- msa/js-executor/queue/kafkaTemplate.js | 43 +++++++++++-------- 4 files changed, 44 insertions(+), 26 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index 867dd82b39..dec935c8ad 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -25,6 +25,7 @@ const config = require('config'), Utils = require('./utils'), JsExecutor = require('./jsExecutor'); +const statFrequency = Number(config.get('script.stat_print_frequency')); const scriptBodyTraceFrequency = Number(config.get('script.script_body_trace_frequency')); const useSandbox = config.get('script.use_sandbox') === 'true'; const maxActiveScripts = Number(config.get('script.max_active_scripts')); @@ -39,6 +40,7 @@ function JsInvokeMessageProcessor(producer) { this.scriptMap = new Map(); this.scriptIds = []; this.executedScriptsCounter = 0; + this.lastStatTime = performance.now(); } JsInvokeMessageProcessor.prototype.onJsInvokeMessage = function(message) { @@ -118,11 +120,16 @@ JsInvokeMessageProcessor.prototype.processInvokeRequest = function(requestId, re var scriptId = getScriptId(invokeRequest); logger.debug('[%s] Processing invoke request, scriptId: [%s]', requestId, scriptId); this.executedScriptsCounter++; - if ( this.executedScriptsCounter >= scriptBodyTraceFrequency ) { - this.executedScriptsCounter = 0; - if (logger.levels[logger.level] >= logger.levels['debug']) { - logger.debug('[%s] Executing script body: [%s]', scriptId, invokeRequest.scriptBody); - } + if (this.executedScriptsCounter % statFrequency == 0) { + var nowMs = performance.now(); + var msSinceLastStat = nowMs - this.lastStatTime; + var requestsPerSec = statFrequency / msSinceLastStat * 1000; //msSinceLastStat can't be zero in the real world + this.lastStatTime = nowMs; + logger.info('STAT[%s]: requests [%s], took [%s]ms, request/s [%s]', this.executedScriptsCounter, statFrequency, msSinceLastStat, requestsPerSec); + } + + if (this.executedScriptsCounter % scriptBodyTraceFrequency == 0) { + logger.info('[%s] Executing script body: [%s]', scriptId, invokeRequest.scriptBody); } this.getOrCompileScript(scriptId, invokeRequest.scriptBody).then( (script) => { diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 5e2bfe8342..c5eb15e8b9 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -26,6 +26,8 @@ kafka: servers: "TB_KAFKA_SERVERS" replication_factor: "TB_QUEUE_KAFKA_REPLICATION_FACTOR" acks: "TB_KAFKA_ACKS" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge + batch_size: "${TB_KAFKA_BATCH_SIZE:128}" # for producer + linger_ms: "${TB_KAFKA_LINGER_MS:1}" # for producer requestTimeout: "TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS" compression: "TB_QUEUE_KAFKA_COMPRESSION" # gzip or uncompressed topic_properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" @@ -70,6 +72,7 @@ logger: script: use_sandbox: "SCRIPT_USE_SANDBOX" + stat_print_frequency: "SCRIPT_STAT_PRINT_FREQUENCY" script_body_trace_frequency: "SCRIPT_BODY_TRACE_FREQUENCY" max_active_scripts: "MAX_ACTIVE_SCRIPTS" slow_query_log_ms: "SLOW_QUERY_LOG_MS" #1.123456 diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index cdf5b35421..5a2c47a675 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -26,6 +26,8 @@ kafka: servers: "localhost:9092" replication_factor: "1" acks: "1" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge + batch_size: "128" # for producer + linger_ms: "1" # for producer requestTimeout: "30000" # The default value in kafkajs is: 30000 compression: "gzip" # gzip or uncompressed topic_properties: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600;partitions:100;min.insync.replicas:1" @@ -59,7 +61,8 @@ logger: script: use_sandbox: "true" - script_body_trace_frequency: "1000" + script_body_trace_frequency: "10000" + stat_print_frequency: "10000" max_active_scripts: "1000" - slow_query_log_ms: "1.000000" #millis + slow_query_log_ms: "5.000000" #millis slow_query_log_body: "false" diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index a061215b28..79d87c17b4 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -23,12 +23,11 @@ const replicationFactor = Number(config.get('kafka.replication_factor')); const topicProperties = config.get('kafka.topic_properties'); const kafkaClientId = config.get('kafka.client_id'); const acks = Number(config.get('kafka.acks')); +const maxBatchSize = Number(config.get('kafka.batch_size')); +const linger = Number(config.get('kafka.linger_ms')); const requestTimeout = Number(config.get('kafka.requestTimeout')); const compressionType = (config.get('kafka.requestTimeout') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; -const linger = 5; //milliseconds //TODO move to the config -const maxBatchSize = 10; //max messages in batch //TODO move to the config - let kafkaClient; let kafkaAdmin; let consumer; @@ -76,7 +75,9 @@ function sendLoopWithLinger() { function sendMessagesAsBatch() { if (batchMessages.length > 0) { - logger.info('sendMessagesAsBatch, lenght: [%s]', batchMessages.length ); + if (batchMessages.length > 1) { + logger.info('sendMessagesAsBatch, length: [%s]', batchMessages.length); + } const messagesToSend = batchMessages; batchMessages = []; producer.sendBatch({ @@ -85,17 +86,15 @@ function sendMessagesAsBatch() { compression: compressionType }).then( () => { - logger.info('Response sent to kafka, length: [%s]', messagesToSend.length ); + logger.debug('Response sent to kafka, length: [%s]', messagesToSend.length); }, (err) => { - logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length ); + logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); batchMessages = messagesToSend.concat(batchMessages); logger.error(err.stack); } ); - } else { - //logger.debug("nothing to send"); } } @@ -163,6 +162,8 @@ function sendMessagesAsBatch() { consumer = kafkaClient.consumer({groupId: 'js-executor-group'}); producer = kafkaClient.producer(); +/* + //producer event instrumentation to debug const { CONNECT } = producer.events; const removeListenerC = producer.on(CONNECT, e => logger.info(`producer CONNECT`)); const { DISCONNECT } = producer.events; @@ -173,18 +174,22 @@ function sendMessagesAsBatch() { const removeListenerRT = producer.on(REQUEST_TIMEOUT, e => logger.info(`producer REQUEST_TIMEOUT ${e.payload.broker}`)); const { REQUEST_QUEUE_SIZE } = producer.events; const removeListenerRQS = producer.on(REQUEST_QUEUE_SIZE, e => logger.info(`producer REQUEST_QUEUE_SIZE ${e.payload.broker} size ${e.queueSize}`)); +*/ -// const removeListeners = {} -// const { FETCH_START } = consumer.events; -// removeListeners[FETCH_START] = consumer.on(FETCH_START, e => logger.info(`consumer FETCH_START`)); -// const { FETCH } = consumer.events; -// removeListeners[FETCH] = consumer.on(FETCH, e => logger.info(`consumer FETCH numberOfBatches ${e.payload.numberOfBatches} duration ${e.payload.duration}`)); -// const { START_BATCH_PROCESS } = consumer.events; -// removeListeners[START_BATCH_PROCESS] = consumer.on(START_BATCH_PROCESS, e => logger.info(`consumer START_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); -// const { END_BATCH_PROCESS } = consumer.events; -// removeListeners[END_BATCH_PROCESS] = consumer.on(END_BATCH_PROCESS, e => logger.info(`consumer END_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); -// const { COMMIT_OFFSETS } = consumer.events; -// removeListeners[COMMIT_OFFSETS] = consumer.on(COMMIT_OFFSETS, e => logger.info(`consumer COMMIT_OFFSETS topics ${e.payload.topics}`)); +/* + //consumer event instrumentation to debug + const removeListeners = {} + const { FETCH_START } = consumer.events; + removeListeners[FETCH_START] = consumer.on(FETCH_START, e => logger.info(`consumer FETCH_START`)); + const { FETCH } = consumer.events; + removeListeners[FETCH] = consumer.on(FETCH, e => logger.info(`consumer FETCH numberOfBatches ${e.payload.numberOfBatches} duration ${e.payload.duration}`)); + const { START_BATCH_PROCESS } = consumer.events; + removeListeners[START_BATCH_PROCESS] = consumer.on(START_BATCH_PROCESS, e => logger.info(`consumer START_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); + const { END_BATCH_PROCESS } = consumer.events; + removeListeners[END_BATCH_PROCESS] = consumer.on(END_BATCH_PROCESS, e => logger.info(`consumer END_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); + const { COMMIT_OFFSETS } = consumer.events; + removeListeners[COMMIT_OFFSETS] = consumer.on(COMMIT_OFFSETS, e => logger.info(`consumer COMMIT_OFFSETS topics ${e.payload.topics}`)); +*/ const messageProcessor = new JsInvokeMessageProcessor(new KafkaProducer()); await consumer.connect(); From 2cba1e2f16e9a1a438d2e4bf7af0afa8968fbae5 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 28 May 2021 18:07:31 +0300 Subject: [PATCH 06/21] js-executor reduced log severity to debug --- msa/js-executor/queue/kafkaTemplate.js | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 79d87c17b4..f467d0813b 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -75,9 +75,7 @@ function sendLoopWithLinger() { function sendMessagesAsBatch() { if (batchMessages.length > 0) { - if (batchMessages.length > 1) { - logger.info('sendMessagesAsBatch, length: [%s]', batchMessages.length); - } + logger.debug('sendMessagesAsBatch, length: [%s]', batchMessages.length); const messagesToSend = batchMessages; batchMessages = []; producer.sendBatch({ From a4e28ad94533610c4de0f855bad75dcd4460ac74 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 1 Jun 2021 16:00:32 +0300 Subject: [PATCH 07/21] js-executor fixed promises for each message for Kafka batches --- .../api/jsInvokeMessageProcessor.js | 24 +++++++------- msa/js-executor/queue/kafkaTemplate.js | 33 ++++++++++++------- 2 files changed, 33 insertions(+), 24 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index dec935c8ad..1aabd216d2 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -180,19 +180,17 @@ JsInvokeMessageProcessor.prototype.sendResponse = function (requestId, responseT var remoteResponse = createRemoteResponse(requestId, compileResponse, invokeResponse, releaseResponse); var rawResponse = Buffer.from(JSON.stringify(remoteResponse), 'utf8'); logger.debug('[%s] Sending response to queue, scriptId: [%s]', requestId, scriptId); - this.producer.send(responseTopic, scriptId, rawResponse, headers); -//TODO put error msg for other queues implementation except Kafka -// .then( -// () => { -// logger.info('[%s] Response sent to queue, took [%s]ms, scriptId: [%s]', requestId, (performance.now() - tStartSending), scriptId); -// }, -// (err) => { -// if (err) { -// logger.error('[%s] Failed to send response to queue: %s', requestId, err.message); -// logger.error(err.stack); -// } -// } -// ); + this.producer.send(responseTopic, scriptId, rawResponse, headers).then( + () => { + logger.debug('[%s] Response sent to queue, took [%s]ms, scriptId: [%s]', requestId, (performance.now() - tStartSending), scriptId); + }, + (err) => { + if (err) { + logger.error('[%s] Failed to send response to queue: %s', requestId, err.message); + logger.error(err.stack); + } + } + ); } JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scriptBody) { diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index f467d0813b..1fac350976 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -36,11 +36,12 @@ let producer; const configEntries = []; let batchMessages = []; +let batchResolvers = []; let sendLoopInstance; function KafkaProducer() { this.send = (responseTopic, scriptId, rawResponse, headers) => { - logger.debug('Pending queue response, scriptId: [%s]', scriptId); + logger.debug('Pending queue response, scriptId: [%s]', scriptId); const message = { topic: responseTopic, messages: [{ @@ -50,17 +51,22 @@ function KafkaProducer() { }] }; - pushMessageToSendLater(message); - return {}; + return pushMessageToSendLater(message); } } function pushMessageToSendLater(message) { + let resolver; + const promise = new Promise((resolve, reject) => { + resolver = resolve; + }); batchMessages.push(message); + batchResolvers.push(resolver); if (batchMessages.length >= maxBatchSize) { sendMessagesAsBatch(); sendLoopWithLinger(); //reset loop function and reschedule new linger } + return promise; } function sendLoopWithLinger() { @@ -77,22 +83,27 @@ function sendMessagesAsBatch() { if (batchMessages.length > 0) { logger.debug('sendMessagesAsBatch, length: [%s]', batchMessages.length); const messagesToSend = batchMessages; + const resolvers = batchResolvers; batchMessages = []; + batchResolvers = []; producer.sendBatch({ topicMessages: messagesToSend, acks: acks, compression: compressionType }).then( - () => { - logger.debug('Response sent to kafka, length: [%s]', messagesToSend.length); - }, - (err) => { - logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); - batchMessages = messagesToSend.concat(batchMessages); - logger.error(err.stack); + () => { + logger.debug('Response batch sent to kafka, length: [%s]', messagesToSend.length); + for (let i = 0; i < promisesToSend.length; i++) { + resolvers[i](); } + }, + (err) => { + logger.error('Failed batch send to kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); + logger.error(err.stack); + batchMessages = messagesToSend.concat(batchMessages); + batchResolvers = resolvers.concat(batchResolvers); //promises will never be rejected. Will retry forever + } ); - } } From 1d9fc4a322f8067b6e35bc96e8ff28eb13f2ca8e Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Tue, 1 Jun 2021 18:01:21 +0300 Subject: [PATCH 08/21] js-executor: typo fix --- msa/js-executor/queue/kafkaTemplate.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 1fac350976..ad243aacdd 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -93,7 +93,7 @@ function sendMessagesAsBatch() { }).then( () => { logger.debug('Response batch sent to kafka, length: [%s]', messagesToSend.length); - for (let i = 0; i < promisesToSend.length; i++) { + for (let i = 0; i < resolvers.length; i++) { resolvers[i](); } }, From 41391dbef8bbfbf1c0f56165dd41277dc45d23fe Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Tue, 1 Jun 2021 18:06:10 +0300 Subject: [PATCH 09/21] js-executor: code format --- .../api/jsInvokeMessageProcessor.js | 56 +++++++++---------- msa/js-executor/queue/kafkaTemplate.js | 46 +++++++-------- 2 files changed, 51 insertions(+), 51 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index 1aabd216d2..bef70936ec 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -43,7 +43,7 @@ function JsInvokeMessageProcessor(producer) { this.lastStatTime = performance.now(); } -JsInvokeMessageProcessor.prototype.onJsInvokeMessage = function(message) { +JsInvokeMessageProcessor.prototype.onJsInvokeMessage = function (message) { var tStart = performance.now(); let requestId; let responseTopic; @@ -78,13 +78,13 @@ JsInvokeMessageProcessor.prototype.onJsInvokeMessage = function(message) { var tFinish = performance.now(); var tTook = tFinish - tStart; - if ( tTook > slowQueryLogMs ) { + if (tTook > slowQueryLogMs) { let functionName; if (request.invokeRequest) { try { buf = Buffer.from(request.invokeRequest['functionName']); functionName = buf.toString('utf8'); - } catch (err){ + } catch (err) { logger.error('[%s] Failed to read functionName from message header: %s', requestId, err.message); logger.error(err.stack); } @@ -97,7 +97,7 @@ JsInvokeMessageProcessor.prototype.onJsInvokeMessage = function(message) { } -JsInvokeMessageProcessor.prototype.processCompileRequest = function(requestId, responseTopic, headers, compileRequest) { +JsInvokeMessageProcessor.prototype.processCompileRequest = function (requestId, responseTopic, headers, compileRequest) { var scriptId = getScriptId(compileRequest); logger.debug('[%s] Processing compile request, scriptId: [%s]', requestId, scriptId); @@ -116,7 +116,7 @@ JsInvokeMessageProcessor.prototype.processCompileRequest = function(requestId, r ); } -JsInvokeMessageProcessor.prototype.processInvokeRequest = function(requestId, responseTopic, headers, invokeRequest) { +JsInvokeMessageProcessor.prototype.processInvokeRequest = function (requestId, responseTopic, headers, invokeRequest) { var scriptId = getScriptId(invokeRequest); logger.debug('[%s] Processing invoke request, scriptId: [%s]', requestId, scriptId); this.executedScriptsCounter++; @@ -160,7 +160,7 @@ JsInvokeMessageProcessor.prototype.processInvokeRequest = function(requestId, re ); } -JsInvokeMessageProcessor.prototype.processReleaseRequest = function(requestId, responseTopic, headers, releaseRequest) { +JsInvokeMessageProcessor.prototype.processReleaseRequest = function (requestId, responseTopic, headers, releaseRequest) { var scriptId = getScriptId(releaseRequest); logger.debug('[%s] Processing release request, scriptId: [%s]', requestId, scriptId); if (this.scriptMap.has(scriptId)) { @@ -193,9 +193,9 @@ JsInvokeMessageProcessor.prototype.sendResponse = function (requestId, responseT ); } -JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scriptBody) { +JsInvokeMessageProcessor.prototype.getOrCompileScript = function (scriptId, scriptBody) { var self = this; - return new Promise(function(resolve, reject) { + return new Promise(function (resolve, reject) { if (self.scriptMap.has(scriptId)) { resolve(self.scriptMap.get(scriptId)); } else { @@ -212,7 +212,7 @@ JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scrip }); } -JsInvokeMessageProcessor.prototype.cacheScript = function(scriptId, script) { +JsInvokeMessageProcessor.prototype.cacheScript = function (scriptId, script) { if (!this.scriptMap.has(scriptId)) { this.scriptIds.push(scriptId); while (this.scriptIds.length > maxActiveScripts) { @@ -229,40 +229,40 @@ JsInvokeMessageProcessor.prototype.cacheScript = function(scriptId, script) { function createRemoteResponse(requestId, compileResponse, invokeResponse, releaseResponse) { const requestIdBits = Utils.UUIDToBits(requestId); return { - requestIdMSB: requestIdBits[0], - requestIdLSB: requestIdBits[1], - compileResponse: compileResponse, - invokeResponse: invokeResponse, - releaseResponse: releaseResponse + requestIdMSB: requestIdBits[0], + requestIdLSB: requestIdBits[1], + compileResponse: compileResponse, + invokeResponse: invokeResponse, + releaseResponse: releaseResponse }; } function createCompileResponse(scriptId, success, errorCode, err) { const scriptIdBits = Utils.UUIDToBits(scriptId); - return { - errorCode: errorCode, - success: success, - errorDetails: parseJsErrorDetails(err), - scriptIdMSB: scriptIdBits[0], - scriptIdLSB: scriptIdBits[1] + return { + errorCode: errorCode, + success: success, + errorDetails: parseJsErrorDetails(err), + scriptIdMSB: scriptIdBits[0], + scriptIdLSB: scriptIdBits[1] }; } function createInvokeResponse(result, success, errorCode, err) { - return { - errorCode: errorCode, - success: success, - errorDetails: parseJsErrorDetails(err), - result: result + return { + errorCode: errorCode, + success: success, + errorDetails: parseJsErrorDetails(err), + result: result }; } function createReleaseResponse(scriptId, success) { const scriptIdBits = Utils.UUIDToBits(scriptId); return { - success: success, - scriptIdMSB: scriptIdBits[0], - scriptIdLSB: scriptIdBits[1] + success: success, + scriptIdMSB: scriptIdBits[0], + scriptIdLSB: scriptIdBits[1] }; } diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index ad243aacdd..7fce429d23 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -45,9 +45,9 @@ function KafkaProducer() { const message = { topic: responseTopic, messages: [{ - key: scriptId, - value: rawResponse, - headers: headers.data + key: scriptId, + value: rawResponse, + headers: headers.data }] }; @@ -70,28 +70,28 @@ function pushMessageToSendLater(message) { } function sendLoopWithLinger() { - if (sendLoopInstance) { + if (sendLoopInstance) { logger.debug("Clear sendLoop scheduler. Starting new send loop with linger [%s]", linger); clearInterval(sendLoopInstance); - } else { + } else { logger.debug("Starting new send loop with linger [%s]", linger) - } - sendLoopInstance = setInterval(sendMessagesAsBatch, linger); + } + sendLoopInstance = setInterval(sendMessagesAsBatch, linger); } function sendMessagesAsBatch() { - if (batchMessages.length > 0) { - logger.debug('sendMessagesAsBatch, length: [%s]', batchMessages.length); - const messagesToSend = batchMessages; - const resolvers = batchResolvers; - batchMessages = []; - batchResolvers = []; - producer.sendBatch({ - topicMessages: messagesToSend, - acks: acks, - compression: compressionType - }).then( - () => { + if (batchMessages.length > 0) { + logger.debug('sendMessagesAsBatch, length: [%s]', batchMessages.length); + const messagesToSend = batchMessages; + const resolvers = batchResolvers; + batchMessages = []; + batchResolvers = []; + producer.sendBatch({ + topicMessages: messagesToSend, + acks: acks, + compression: compressionType + }).then( + () => { logger.debug('Response batch sent to kafka, length: [%s]', messagesToSend.length); for (let i = 0; i < resolvers.length; i++) { resolvers[i](); @@ -103,8 +103,8 @@ function sendMessagesAsBatch() { batchMessages = messagesToSend.concat(batchMessages); batchResolvers = resolvers.concat(batchResolvers); //promises will never be rejected. Will retry forever } - ); - } + ); + } } (async () => { @@ -120,8 +120,8 @@ function sendMessagesAsBatch() { let kafkaConfig = { brokers: kafkaBootstrapServers.split(','), - logLevel: logLevel.INFO, - logCreator: KafkaJsWinstonLogCreator + logLevel: logLevel.INFO, + logCreator: KafkaJsWinstonLogCreator }; if (kafkaClientId) { From ab5f1b5b639236fb7528ff9544425e2310c33dbb Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Wed, 2 Jun 2021 10:41:24 +0300 Subject: [PATCH 10/21] js-executor: ScriptMap optimize work --- msa/js-executor/api/jsInvokeMessageProcessor.js | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index bef70936ec..aa8f5d2db0 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -196,8 +196,9 @@ JsInvokeMessageProcessor.prototype.sendResponse = function (requestId, responseT JsInvokeMessageProcessor.prototype.getOrCompileScript = function (scriptId, scriptBody) { var self = this; return new Promise(function (resolve, reject) { - if (self.scriptMap.has(scriptId)) { - resolve(self.scriptMap.get(scriptId)); + const script = self.scriptMap.get(scriptId); + if (script !== undefined) { + resolve(script); } else { self.executor.compileScript(scriptBody).then( (script) => { From ff7fa6237f1f342357995c32488f542c984583d9 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 2 Jun 2021 11:05:32 +0300 Subject: [PATCH 11/21] js-executor: zero check and code cleanup --- msa/js-executor/api/jsInvokeMessageProcessor.js | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index aa8f5d2db0..5257a03168 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -121,9 +121,9 @@ JsInvokeMessageProcessor.prototype.processInvokeRequest = function (requestId, r logger.debug('[%s] Processing invoke request, scriptId: [%s]', requestId, scriptId); this.executedScriptsCounter++; if (this.executedScriptsCounter % statFrequency == 0) { - var nowMs = performance.now(); - var msSinceLastStat = nowMs - this.lastStatTime; - var requestsPerSec = statFrequency / msSinceLastStat * 1000; //msSinceLastStat can't be zero in the real world + const nowMs = performance.now(); + const msSinceLastStat = nowMs - this.lastStatTime; + const requestsPerSec = msSinceLastStat == 0 ? statFrequency : statFrequency / msSinceLastStat * 1000; this.lastStatTime = nowMs; logger.info('STAT[%s]: requests [%s], took [%s]ms, request/s [%s]', this.executedScriptsCounter, statFrequency, msSinceLastStat, requestsPerSec); } @@ -197,13 +197,13 @@ JsInvokeMessageProcessor.prototype.getOrCompileScript = function (scriptId, scri var self = this; return new Promise(function (resolve, reject) { const script = self.scriptMap.get(scriptId); - if (script !== undefined) { + if (script) { resolve(script); } else { self.executor.compileScript(scriptBody).then( - (script) => { - self.cacheScript(scriptId, script); - resolve(script); + (compiledScript) => { + self.cacheScript(scriptId, compiledScript); + resolve(compiledScript); }, (err) => { reject(err); From dd748245706903847ba72f7c6066e97e026ebe28 Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Wed, 9 Jun 2021 13:30:28 +0300 Subject: [PATCH 12/21] js-executor: Refactoring Kafka executor for use async/await --- msa/js-executor/queue/kafkaTemplate.js | 68 +++++++++++--------------- 1 file changed, 28 insertions(+), 40 deletions(-) diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 7fce429d23..b34727c977 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -36,11 +36,10 @@ let producer; const configEntries = []; let batchMessages = []; -let batchResolvers = []; let sendLoopInstance; function KafkaProducer() { - this.send = (responseTopic, scriptId, rawResponse, headers) => { + this.send = async (responseTopic, scriptId, rawResponse, headers) => { logger.debug('Pending queue response, scriptId: [%s]', scriptId); const message = { topic: responseTopic, @@ -51,60 +50,50 @@ function KafkaProducer() { }] }; - return pushMessageToSendLater(message); + await pushMessageToSendLater(message); } } -function pushMessageToSendLater(message) { - let resolver; - const promise = new Promise((resolve, reject) => { - resolver = resolve; - }); +async function pushMessageToSendLater(message) { batchMessages.push(message); - batchResolvers.push(resolver); if (batchMessages.length >= maxBatchSize) { - sendMessagesAsBatch(); - sendLoopWithLinger(); //reset loop function and reschedule new linger + await sendMessagesAsBatch(true); } - return promise; } function sendLoopWithLinger() { if (sendLoopInstance) { - logger.debug("Clear sendLoop scheduler. Starting new send loop with linger [%s]", linger); - clearInterval(sendLoopInstance); + clearTimeout(sendLoopInstance); } else { logger.debug("Starting new send loop with linger [%s]", linger) } - sendLoopInstance = setInterval(sendMessagesAsBatch, linger); + sendLoopInstance = setTimeout(sendMessagesAsBatch, linger); } -function sendMessagesAsBatch() { +async function sendMessagesAsBatch(isImmediately) { + if (sendLoopInstance) { + logger.debug("sendMessagesAsBatch: Clear sendLoop scheduler. Starting new send loop with linger [%s]", linger); + clearTimeout(sendLoopInstance); + } + sendLoopInstance = null; if (batchMessages.length > 0) { - logger.debug('sendMessagesAsBatch, length: [%s]', batchMessages.length); + logger.debug('sendMessagesAsBatch, length: [%s], %s', batchMessages.length, isImmediately ? 'immediately' : ''); const messagesToSend = batchMessages; - const resolvers = batchResolvers; batchMessages = []; - batchResolvers = []; - producer.sendBatch({ - topicMessages: messagesToSend, - acks: acks, - compression: compressionType - }).then( - () => { - logger.debug('Response batch sent to kafka, length: [%s]', messagesToSend.length); - for (let i = 0; i < resolvers.length; i++) { - resolvers[i](); - } - }, - (err) => { - logger.error('Failed batch send to kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); - logger.error(err.stack); - batchMessages = messagesToSend.concat(batchMessages); - batchResolvers = resolvers.concat(batchResolvers); //promises will never be rejected. Will retry forever - } - ); + try { + await producer.sendBatch({ + topicMessages: messagesToSend, + acks: acks, + compression: compressionType + }) + logger.debug('Response batch sent to kafka, length: [%s]', messagesToSend.length); + } catch(err) { + logger.error('Failed batch send to kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); + logger.error(err.stack); + batchMessages = messagesToSend.concat(batchMessages); + } } + sendLoopWithLinger(); } (async () => { @@ -285,9 +274,8 @@ async function disconnectProducer() { producer = null; try { logger.info('Stopping loop...'); - //TODO: send handle msg - clearInterval(sendLoopInstance); - sendMessagesAsBatch(); + clearTimeout(sendLoopInstance); + await sendMessagesAsBatch(); await _producer.disconnect(); logger.info('Kafka Producer stopped.'); } catch (e) { From a7a7dca3b3eb76166cf8ed93bf764f2c8c5935b2 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 8 Jun 2021 20:03:45 +0300 Subject: [PATCH 13/21] js-executor: upgraded libraries to fix vulnerability warning (@google-cloud/pubsub", amqplib) --- msa/js-executor/package.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index 75f009880b..280979afa6 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -13,8 +13,8 @@ }, "dependencies": { "@azure/service-bus": "^1.1.9", - "@google-cloud/pubsub": "^2.5.0", - "amqplib": "^0.6.0", + "@google-cloud/pubsub": "^2.12.0", + "amqplib": "^0.8.0", "aws-sdk": "^2.741.0", "azure-sb": "^0.11.1", "config": "^3.3.1", From 98bea8eaf579f6641bc8a3dfa43639901aca7c91 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 8 Jun 2021 20:14:10 +0300 Subject: [PATCH 14/21] js-executor: http livenessProbe added --- msa/js-executor/api/httpServer.js | 46 +++++++++++++++++++ .../config/custom-environment-variables.yml | 1 + msa/js-executor/config/default.yml | 1 + msa/js-executor/package.json | 1 + msa/js-executor/server.js | 2 + 5 files changed, 51 insertions(+) create mode 100644 msa/js-executor/api/httpServer.js diff --git a/msa/js-executor/api/httpServer.js b/msa/js-executor/api/httpServer.js new file mode 100644 index 0000000000..996638c7c3 --- /dev/null +++ b/msa/js-executor/api/httpServer.js @@ -0,0 +1,46 @@ +/* + * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * + * Copyright © 2016-2021 ThingsBoard, Inc. All Rights Reserved. + * + * NOTICE: All information contained herein is, and remains + * the property of ThingsBoard, Inc. and its suppliers, + * if any. The intellectual and technical concepts contained + * herein are proprietary to ThingsBoard, Inc. + * and its suppliers and may be covered by U.S. and Foreign Patents, + * patents in process, and are protected by trade secret or copyright law. + * + * Dissemination of this information or reproduction of this material is strictly forbidden + * unless prior written permission is obtained from COMPANY. + * + * Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees, + * managers or contractors who have executed Confidentiality and Non-disclosure agreements + * explicitly covering such access. + * + * The copyright notice above does not evidence any actual or intended publication + * or disclosure of this source code, which includes + * information that is confidential and/or proprietary, and is a trade secret, of COMPANY. + * ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE, + * OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT + * THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED, + * AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES. + * THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION + * DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS, + * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. + */ + +const config = require('config'), + logger = require('../config/logger')._logger('httpServer'), + express = require('express'); + +const httpPort = Number(config.get('http_port')); + +const app = express(); + +app.get('/livenessProbe', async (req, res) => { + const date = new Date(); + const message = { now: date.toISOString() }; + res.send(message); +}) + +app.listen(httpPort, () => logger.info(`Started http endpoint on port ${httpPort}. Please, use /livenessProbe !`)) \ No newline at end of file diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index c5eb15e8b9..7773fdaf0c 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -16,6 +16,7 @@ queue_type: "TB_QUEUE_TYPE" #kafka (Apache Kafka) or aws-sqs (AWS SQS) or pubsub (PubSub) or service-bus (Azure Service Bus) or rabbitmq (RabbitMQ) request_topic: "REMOTE_JS_EVAL_REQUEST_TOPIC" +http_port: "HTTP_PORT" # /livenessProbe js: response_poll_interval: "REMOTE_JS_RESPONSE_POLL_INTERVAL_MS" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index 5a2c47a675..bcefc1c639 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -16,6 +16,7 @@ queue_type: "kafka" request_topic: "js_eval.requests" +http_port: "8888" # /livenessProbe js: response_poll_interval: "25" diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index 280979afa6..ba940ee85a 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -18,6 +18,7 @@ "aws-sdk": "^2.741.0", "azure-sb": "^0.11.1", "config": "^3.3.1", + "express": "^4.17.1", "js-yaml": "^3.14.0", "kafkajs": "^1.15.0", "long": "^4.0.0", diff --git a/msa/js-executor/server.js b/msa/js-executor/server.js index 0c415cc156..7b2e7e59b8 100644 --- a/msa/js-executor/server.js +++ b/msa/js-executor/server.js @@ -51,3 +51,5 @@ switch (serviceType) { process.exit(-1); } +require('./api/httpServer'); + From 05df5e03efbaff6a1a5b48d27dc6327962ecd4b2 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 10 Jun 2021 11:26:33 +0300 Subject: [PATCH 15/21] Revert "js-executor: upgraded libraries to fix vulnerability warning (@google-cloud/pubsub", amqplib)" Untested version with breaking changes declared This reverts commit 80f83a86 --- msa/js-executor/package.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index ba940ee85a..187144da20 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -13,8 +13,8 @@ }, "dependencies": { "@azure/service-bus": "^1.1.9", - "@google-cloud/pubsub": "^2.12.0", - "amqplib": "^0.8.0", + "@google-cloud/pubsub": "^2.5.0", + "amqplib": "^0.6.0", "aws-sdk": "^2.741.0", "azure-sb": "^0.11.1", "config": "^3.3.1", From 3a31b9c5ea7451dce8e5936e196f9fdbbb75eb4c Mon Sep 17 00:00:00 2001 From: Kien Truong Date: Thu, 3 Jun 2021 11:05:02 +0700 Subject: [PATCH 16/21] Fix wrong configuration key for Compression Type Fix #4678 (cherry picked from commit c176ca94aa0f612c3bb3856ae2acf3a399cc74a4) --- msa/js-executor/queue/kafkaTemplate.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index b34727c977..7e19f373a8 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -26,7 +26,7 @@ const acks = Number(config.get('kafka.acks')); const maxBatchSize = Number(config.get('kafka.batch_size')); const linger = Number(config.get('kafka.linger_ms')); const requestTimeout = Number(config.get('kafka.requestTimeout')); -const compressionType = (config.get('kafka.requestTimeout') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; +const compressionType = (config.get('kafka.compression') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; let kafkaClient; let kafkaAdmin; From 9a4a94621a88dd7e6e329dc1d7767ff75e311e16 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 10 Jun 2021 14:31:51 +0300 Subject: [PATCH 17/21] js-executor: parameter added for Kafka PARTITIONS_CONSUMED_CONCURRENTLY to decrease max latency while scale down replicas --- msa/js-executor/config/custom-environment-variables.yml | 1 + msa/js-executor/config/default.yml | 1 + msa/js-executor/queue/kafkaTemplate.js | 3 ++- 3 files changed, 4 insertions(+), 1 deletion(-) diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 7773fdaf0c..34903c4541 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -29,6 +29,7 @@ kafka: acks: "TB_KAFKA_ACKS" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge batch_size: "${TB_KAFKA_BATCH_SIZE:128}" # for producer linger_ms: "${TB_KAFKA_LINGER_MS:1}" # for producer + partitions_consumed_concurrently: "${PARTITIONS_CONSUMED_CONCURRENTLY:1}" # increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency requestTimeout: "TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS" compression: "TB_QUEUE_KAFKA_COMPRESSION" # gzip or uncompressed topic_properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index bcefc1c639..190cf097ed 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -29,6 +29,7 @@ kafka: acks: "1" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge batch_size: "128" # for producer linger_ms: "1" # for producer + partitions_consumed_concurrently: "1" # increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency requestTimeout: "30000" # The default value in kafkajs is: 30000 compression: "gzip" # gzip or uncompressed topic_properties: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600;partitions:100;min.insync.replicas:1" diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 7e19f373a8..46ded135e6 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -27,6 +27,7 @@ const maxBatchSize = Number(config.get('kafka.batch_size')); const linger = Number(config.get('kafka.linger_ms')); const requestTimeout = Number(config.get('kafka.requestTimeout')); const compressionType = (config.get('kafka.compression') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; +const partitionsConsumedConcurrently = Number(config.get('kafka.partitions_consumed_concurrently')); let kafkaClient; let kafkaAdmin; @@ -197,7 +198,7 @@ async function sendMessagesAsBatch(isImmediately) { logger.info('Started ThingsBoard JavaScript Executor Microservice.'); await consumer.run({ - //partitionsConsumedConcurrently: 1, // Default: 1 + partitionsConsumedConcurrently: partitionsConsumedConcurrently, eachMessage: async ({topic, partition, message}) => { let headers = message.headers; let key = message.key; From 9ffa50f4a528d3f7b99e167e3bd797213ed26314 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 10 Jun 2021 16:05:50 +0300 Subject: [PATCH 18/21] js-executor: fixed env variables mapping --- msa/js-executor/config/custom-environment-variables.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 34903c4541..4185e65835 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -27,9 +27,9 @@ kafka: servers: "TB_KAFKA_SERVERS" replication_factor: "TB_QUEUE_KAFKA_REPLICATION_FACTOR" acks: "TB_KAFKA_ACKS" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge - batch_size: "${TB_KAFKA_BATCH_SIZE:128}" # for producer - linger_ms: "${TB_KAFKA_LINGER_MS:1}" # for producer - partitions_consumed_concurrently: "${PARTITIONS_CONSUMED_CONCURRENTLY:1}" # increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency + batch_size: "TB_KAFKA_BATCH_SIZE" # for producer + linger_ms: "TB_KAFKA_LINGER_MS" # for producer + partitions_consumed_concurrently: "TB_KAFKA_PARTITIONS_CONSUMED_CONCURRENTLY" # increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency requestTimeout: "TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS" compression: "TB_QUEUE_KAFKA_COMPRESSION" # gzip or uncompressed topic_properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" From cda6f3678b8cfe7a0046d75ff5daeae3287b8312 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 07:27:42 +0300 Subject: [PATCH 19/21] js-executor: fixed license header --- msa/js-executor/api/httpServer.js | 36 +++++++++---------------------- 1 file changed, 10 insertions(+), 26 deletions(-) diff --git a/msa/js-executor/api/httpServer.js b/msa/js-executor/api/httpServer.js index 996638c7c3..4948c45f61 100644 --- a/msa/js-executor/api/httpServer.js +++ b/msa/js-executor/api/httpServer.js @@ -1,34 +1,18 @@ /* - * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * Copyright © 2016-2021 The Thingsboard Authors * - * Copyright © 2016-2021 ThingsBoard, Inc. All Rights Reserved. + * 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 * - * NOTICE: All information contained herein is, and remains - * the property of ThingsBoard, Inc. and its suppliers, - * if any. The intellectual and technical concepts contained - * herein are proprietary to ThingsBoard, Inc. - * and its suppliers and may be covered by U.S. and Foreign Patents, - * patents in process, and are protected by trade secret or copyright law. + * http://www.apache.org/licenses/LICENSE-2.0 * - * Dissemination of this information or reproduction of this material is strictly forbidden - * unless prior written permission is obtained from COMPANY. - * - * Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees, - * managers or contractors who have executed Confidentiality and Non-disclosure agreements - * explicitly covering such access. - * - * The copyright notice above does not evidence any actual or intended publication - * or disclosure of this source code, which includes - * information that is confidential and/or proprietary, and is a trade secret, of COMPANY. - * ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE, - * OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT - * THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED, - * AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES. - * THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION - * DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS, - * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. + * 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. */ - const config = require('config'), logger = require('../config/logger')._logger('httpServer'), express = require('express'); From c07379812b213f3ca249da31285a00124d6345bd Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 07:28:53 +0300 Subject: [PATCH 20/21] js-executor: updated yarn.lock --- msa/js-executor/yarn.lock | 361 ++++++++++++++++++++++++++++++++++++-- 1 file changed, 347 insertions(+), 14 deletions(-) diff --git a/msa/js-executor/yarn.lock b/msa/js-executor/yarn.lock index 404c059c23..0092372b69 100644 --- a/msa/js-executor/yarn.lock +++ b/msa/js-executor/yarn.lock @@ -418,6 +418,14 @@ abort-controller@^3.0.0: dependencies: event-target-shim "^5.0.0" +accepts@~1.3.7: + version "1.3.7" + resolved "https://registry.yarnpkg.com/accepts/-/accepts-1.3.7.tgz#531bc726517a3b2b41f850021c6cc15eaab507cd" + integrity sha512-Il80Qs2WjYlJIBNzNkK6KYqlVMTbZLXgHx2oT0pU/fjRHyEp+PEfEPY0R3WCwAGVOtauxh1hOxNgIf5bv7dQpA== + dependencies: + mime-types "~2.1.24" + negotiator "0.6.2" + agent-base@6: version "6.0.1" resolved "https://registry.yarnpkg.com/agent-base/-/agent-base-6.0.1.tgz#808007e4e5867decb0ab6ab2f928fbdb5a596db4" @@ -487,6 +495,11 @@ argparse@^1.0.7: dependencies: sprintf-js "~1.0.2" +array-flatten@1.1.1: + version "1.1.1" + resolved "https://registry.yarnpkg.com/array-flatten/-/array-flatten-1.1.1.tgz#9a5f699051b1e7073328f2a008968b64ea2955d2" + integrity sha1-ml9pkFGx5wczKPKgCJaLZOopVdI= + array-union@^2.1.0: version "2.1.0" resolved "https://registry.yarnpkg.com/array-union/-/array-union-2.1.0.tgz#b798420adbeb1de828d84acd8a2e23d3efe85e8d" @@ -621,6 +634,22 @@ bluebird@^3.5.2: resolved "https://registry.yarnpkg.com/bluebird/-/bluebird-3.7.2.tgz#9f229c15be272454ffa973ace0dbee79a1b0c36f" integrity sha512-XpNj6GDQzdfW+r2Wnn7xiSAd7TM3jzkxGXBGTtWKuSXv1xUV+azxAm8jdWZN06QTQk+2N2XB9jRDkvbmQmcRtg== +body-parser@1.19.0: + version "1.19.0" + resolved "https://registry.yarnpkg.com/body-parser/-/body-parser-1.19.0.tgz#96b2709e57c9c4e09a6fd66a8fd979844f69f08a" + integrity sha512-dhEPs72UPbDnAQJ9ZKMNTP6ptJaionhP5cBb541nXPlW60Jepo9RV/a4fX4XWW9CuFNK22krhrj1+rgzifNCsw== + dependencies: + bytes "3.1.0" + content-type "~1.0.4" + debug "2.6.9" + depd "~1.1.2" + http-errors "1.7.2" + iconv-lite "0.4.24" + on-finished "~2.3.0" + qs "6.7.0" + raw-body "2.4.0" + type-is "~1.6.17" + boxen@^4.2.0: version "4.2.0" resolved "https://registry.yarnpkg.com/boxen/-/boxen-4.2.0.tgz#e411b62357d6d6d36587c8ac3d5d974daa070e64" @@ -682,6 +711,11 @@ byline@^5.0.0: resolved "https://registry.yarnpkg.com/byline/-/byline-5.0.0.tgz#741c5216468eadc457b03410118ad77de8c1ddb1" integrity sha1-dBxSFkaOrcRXsDQQEYrXfejB3bE= +bytes@3.1.0: + version "3.1.0" + resolved "https://registry.yarnpkg.com/bytes/-/bytes-3.1.0.tgz#f6cf7933a360e0588fa9fde85651cdc7f805d1f6" + integrity sha512-zauLjrfCG+xvoyaqLoV8bLVXXNGC4JqlxFCutSDWA6fJrTo2ZuvLYTqZ7aHBLZSMOopbzwv8f+wZcVzfVTI2Dg== + cacheable-request@^6.0.0: version "6.1.0" resolved "https://registry.yarnpkg.com/cacheable-request/-/cacheable-request-6.1.0.tgz#20ffb8bd162ba4be11e9567d823db651052ca912" @@ -838,6 +872,28 @@ configstore@^5.0.1: write-file-atomic "^3.0.0" xdg-basedir "^4.0.0" +content-disposition@0.5.3: + version "0.5.3" + resolved "https://registry.yarnpkg.com/content-disposition/-/content-disposition-0.5.3.tgz#e130caf7e7279087c5616c2007d0485698984fbd" + integrity sha512-ExO0774ikEObIAEV9kDo50o+79VCUdEB6n6lzKgGwupcVeRlhrj3qGAfwq8G6uBJjkqLrhT0qEYFcWng8z1z0g== + dependencies: + safe-buffer "5.1.2" + +content-type@~1.0.4: + version "1.0.4" + resolved "https://registry.yarnpkg.com/content-type/-/content-type-1.0.4.tgz#e138cc75e040c727b1966fe5e5f8c9aee256fe3b" + integrity sha512-hIP3EEPs8tB9AT1L+NUqtwOAps4mk2Zob89MWXMHjHWg9milF/j4osnnQLXBCBFBk/tvIG/tUc9mOUJiPBhPXA== + +cookie-signature@1.0.6: + version "1.0.6" + resolved "https://registry.yarnpkg.com/cookie-signature/-/cookie-signature-1.0.6.tgz#e303a882b342cc3ee8ca513a79999734dab3ae2c" + integrity sha1-4wOogrNCzD7oylE6eZmXNNqzriw= + +cookie@0.4.0: + version "0.4.0" + resolved "https://registry.yarnpkg.com/cookie/-/cookie-0.4.0.tgz#beb437e7022b3b6d49019d088665303ebe9c14ba" + integrity sha512-+Hp8fLp57wnUSt0tY0tHEXh4voZRDnoIrZPqlo3DPiI4y9lwg/jqx+1Om94/W6ZaPDOUbnjOt/99w66zk+l1Xg== + core-util-is@1.0.2, core-util-is@~1.0.0: version "1.0.2" resolved "https://registry.yarnpkg.com/core-util-is/-/core-util-is-1.0.2.tgz#b5fd54220aa2bc5ab57aab7140c940754503c1a7" @@ -867,6 +923,13 @@ dateformat@1.0.2-1.2.3: dependencies: ms "^2.1.1" +debug@2.6.9, debug@^2.2.0, debug@~2.6.9: + version "2.6.9" + resolved "https://registry.yarnpkg.com/debug/-/debug-2.6.9.tgz#5d128515df134ff327e90a4c93f4e077a536341f" + integrity sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA== + dependencies: + ms "2.0.0" + debug@4, debug@^4.1.1: version "4.1.1" resolved "https://registry.yarnpkg.com/debug/-/debug-4.1.1.tgz#3b72260255109c6b589cee050f1d516139664791" @@ -874,13 +937,6 @@ debug@4, debug@^4.1.1: dependencies: ms "^2.1.1" -debug@^2.2.0, debug@~2.6.9: - version "2.6.9" - resolved "https://registry.yarnpkg.com/debug/-/debug-2.6.9.tgz#5d128515df134ff327e90a4c93f4e077a536341f" - integrity sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA== - dependencies: - ms "2.0.0" - decamelize@^1.2.0: version "1.2.0" resolved "https://registry.yarnpkg.com/decamelize/-/decamelize-1.2.0.tgz#f6534d15148269b20352e7bee26f501f9a191290" @@ -913,6 +969,16 @@ delayed-stream@~1.0.0: resolved "https://registry.yarnpkg.com/delayed-stream/-/delayed-stream-1.0.0.tgz#df3ae199acadfb7d440aaae0b29e2272b24ec619" integrity sha1-3zrhmayt+31ECqrgsp4icrJOxhk= +depd@~1.1.2: + version "1.1.2" + resolved "https://registry.yarnpkg.com/depd/-/depd-1.1.2.tgz#9bcd52e14c097763e749b274c4346ed2e560b5a9" + integrity sha1-m81S4UwJd2PnSbJ0xDRu0uVgtak= + +destroy@~1.0.4: + version "1.0.4" + resolved "https://registry.yarnpkg.com/destroy/-/destroy-1.0.4.tgz#978857442c44749e4206613e37946205826abd80" + integrity sha1-l4hXRCxEdJ5CBmE+N5RiBYJqvYA= + dir-glob@^3.0.1: version "3.0.1" resolved "https://registry.yarnpkg.com/dir-glob/-/dir-glob-3.0.1.tgz#56dbf73d992a4a93ba1584f4534063fd2e41717f" @@ -962,6 +1028,11 @@ ecdsa-sig-formatter@1.0.11, ecdsa-sig-formatter@^1.0.11: dependencies: safe-buffer "^5.0.1" +ee-first@1.1.1: + version "1.1.1" + resolved "https://registry.yarnpkg.com/ee-first/-/ee-first-1.1.1.tgz#590c61156b0ae2f4f0255732a158b266bc56b21d" + integrity sha1-WQxhFWsK4vTwJVcyoViyZrxWsh0= + emoji-regex@^7.0.1: version "7.0.3" resolved "https://registry.yarnpkg.com/emoji-regex/-/emoji-regex-7.0.3.tgz#933a04052860c85e83c122479c4748a8e4c72156" @@ -977,6 +1048,11 @@ enabled@2.0.x: resolved "https://registry.yarnpkg.com/enabled/-/enabled-2.0.0.tgz#f9dd92ec2d6f4bbc0d5d1e64e21d61cd4665e7c2" integrity sha512-AKrN98kuwOzMIdAizXGI86UFBoo26CL21UM763y1h/GMSJ4/OHU9k2YlsmBpyScFo/wbLzWQJBMCW4+IO3/+OQ== +encodeurl@~1.0.2: + version "1.0.2" + resolved "https://registry.yarnpkg.com/encodeurl/-/encodeurl-1.0.2.tgz#ad3ff4c86ec2d029322f5a02c3a9a606c95b3f59" + integrity sha1-rT/0yG7C0CkyL1oCw6mmBslbP1k= + end-of-stream@^1.0.0, end-of-stream@^1.1.0: version "1.4.4" resolved "https://registry.yarnpkg.com/end-of-stream/-/end-of-stream-1.4.4.tgz#5ae64a5f45057baf3626ec14da0ca5e4b2431eb0" @@ -994,6 +1070,11 @@ escape-goat@^2.0.0: resolved "https://registry.yarnpkg.com/escape-goat/-/escape-goat-2.1.1.tgz#1b2dc77003676c457ec760b2dc68edb648188675" integrity sha512-8/uIhbG12Csjy2JEW7D9pHbreaVaS/OpN3ycnyvElTdwM5n6GY6W6e2IPemfvGZeUMqZ9A/3GqIZMgKnBhAw/Q== +escape-html@~1.0.3: + version "1.0.3" + resolved "https://registry.yarnpkg.com/escape-html/-/escape-html-1.0.3.tgz#0258eae4d3d0c0974de1c169188ef0051d1d1988" + integrity sha1-Aljq5NPQwJdN4cFpGI7wBR0dGYg= + escodegen@^1.14.1: version "1.14.3" resolved "https://registry.yarnpkg.com/escodegen/-/escodegen-1.14.3.tgz#4e7b81fba61581dc97582ed78cab7f0e8d63f503" @@ -1021,6 +1102,11 @@ esutils@^2.0.2: resolved "https://registry.yarnpkg.com/esutils/-/esutils-2.0.3.tgz#74d2eb4de0b8da1293711910d50775b9b710ef64" integrity sha512-kVscqXk4OCp68SZ0dkgEKVi6/8ij300KBWTJq32P/dYeWTSwK41WyTxalN1eRmA5Z9UU/LX9D7FWSmV9SAYx6g== +etag@~1.8.1: + version "1.8.1" + resolved "https://registry.yarnpkg.com/etag/-/etag-1.8.1.tgz#41ae2eeb65efa62268aebfea83ac7d79299b0887" + integrity sha1-Qa4u62XvpiJorr/qg6x9eSmbCIc= + event-target-shim@^5.0.0: version "5.0.1" resolved "https://registry.yarnpkg.com/event-target-shim/-/event-target-shim-5.0.1.tgz#5d4d3ebdf9583d63a5333ce2deb7480ab2b05789" @@ -1041,6 +1127,42 @@ expand-template@^2.0.3: resolved "https://registry.yarnpkg.com/expand-template/-/expand-template-2.0.3.tgz#6e14b3fcee0f3a6340ecb57d2e8918692052a47c" integrity sha512-XYfuKMvj4O35f/pOXLObndIRvyQ+/+6AhODh+OKWj9S9498pHHn/IMszH+gt0fBCRWMNfk1ZSp5x3AifmnI2vg== +express@^4.17.1: + version "4.17.1" + resolved "https://registry.yarnpkg.com/express/-/express-4.17.1.tgz#4491fc38605cf51f8629d39c2b5d026f98a4c134" + integrity sha512-mHJ9O79RqluphRrcw2X/GTh3k9tVv8YcoyY4Kkh4WDMUYKRZUq0h1o0w2rrrxBqM7VoeUVqgb27xlEMXTnYt4g== + dependencies: + accepts "~1.3.7" + array-flatten "1.1.1" + body-parser "1.19.0" + content-disposition "0.5.3" + content-type "~1.0.4" + cookie "0.4.0" + cookie-signature "1.0.6" + debug "2.6.9" + depd "~1.1.2" + encodeurl "~1.0.2" + escape-html "~1.0.3" + etag "~1.8.1" + finalhandler "~1.1.2" + fresh "0.5.2" + merge-descriptors "1.0.1" + methods "~1.1.2" + on-finished "~2.3.0" + parseurl "~1.3.3" + path-to-regexp "0.1.7" + proxy-addr "~2.0.5" + qs "6.7.0" + range-parser "~1.2.1" + safe-buffer "5.1.2" + send "0.17.1" + serve-static "1.14.1" + setprototypeof "1.1.1" + statuses "~1.5.0" + type-is "~1.6.18" + utils-merge "1.0.1" + vary "~1.1.2" + extend@^3.0.2, extend@~3.0.2: version "3.0.2" resolved "https://registry.yarnpkg.com/extend/-/extend-3.0.2.tgz#f8b1136b4071fbd8eb140aff858b1019ec2915fa" @@ -1119,6 +1241,19 @@ fill-range@^7.0.1: dependencies: to-regex-range "^5.0.1" +finalhandler@~1.1.2: + version "1.1.2" + resolved "https://registry.yarnpkg.com/finalhandler/-/finalhandler-1.1.2.tgz#b7e7d000ffd11938d0fdb053506f6ebabe9f587d" + integrity sha512-aAWcW57uxVNrQZqFXjITpW3sIUQmHGG3qSb9mUah9MgMC4NeWhNOlNjXEYq3HjRAvL6arUviZGGJsBg6z0zsWA== + dependencies: + debug "2.6.9" + encodeurl "~1.0.2" + escape-html "~1.0.3" + on-finished "~2.3.0" + parseurl "~1.3.3" + statuses "~1.5.0" + unpipe "~1.0.0" + find-up@^4.1.0: version "4.1.0" resolved "https://registry.yarnpkg.com/find-up/-/find-up-4.1.0.tgz#97afe7d6cdc0bc5928584b7c8d7b16e8a9aa5d19" @@ -1155,6 +1290,16 @@ form-data@~2.3.2: combined-stream "^1.0.6" mime-types "^2.1.12" +forwarded@0.2.0: + version "0.2.0" + resolved "https://registry.yarnpkg.com/forwarded/-/forwarded-0.2.0.tgz#2269936428aad4c15c7ebe9779a84bf0b2a81811" + integrity sha512-buRG0fpBtRHSTCOASe6hD258tEubFoRLb4ZNA6NxMVHNw2gOcwHo9wyablzMzOA5z9xA9L1KNjk/Nt6MT9aYow== + +fresh@0.5.2: + version "0.5.2" + resolved "https://registry.yarnpkg.com/fresh/-/fresh-0.5.2.tgz#3d8cadd90d976569fa835ab1f8e4b23a105605a7" + integrity sha1-PYyt2Q2XZWn6g1qx+OSyOhBWBac= + from2@^2.3.0: version "2.3.0" resolved "https://registry.yarnpkg.com/from2/-/from2-2.3.0.tgz#8bfb5502bde4a4d36cfdeea007fcca21d7e382af" @@ -1384,6 +1529,28 @@ http-cache-semantics@^4.0.0: resolved "https://registry.yarnpkg.com/http-cache-semantics/-/http-cache-semantics-4.1.0.tgz#49e91c5cbf36c9b94bcfcd71c23d5249ec74e390" integrity sha512-carPklcUh7ROWRK7Cv27RPtdhYhUsela/ue5/jKzjegVvXDqM2ILE9Q2BGn9JZJh1g87cp56su/FgQSzcWS8cQ== +http-errors@1.7.2: + version "1.7.2" + resolved "https://registry.yarnpkg.com/http-errors/-/http-errors-1.7.2.tgz#4f5029cf13239f31036e5b2e55292bcfbcc85c8f" + integrity sha512-uUQBt3H/cSIVfch6i1EuPNy/YsRSOUBXTVfZ+yR7Zjez3qjBz6i9+i4zjNaoqcoFVI4lQJ5plg63TvGfRSDCRg== + dependencies: + depd "~1.1.2" + inherits "2.0.3" + setprototypeof "1.1.1" + statuses ">= 1.5.0 < 2" + toidentifier "1.0.0" + +http-errors@~1.7.2: + version "1.7.3" + resolved "https://registry.yarnpkg.com/http-errors/-/http-errors-1.7.3.tgz#6c619e4f9c60308c38519498c14fbb10aacebb06" + integrity sha512-ZTTX0MWrsQ2ZAhA1cejAwDLycFsd7I7nVtnkT3Ol0aqodaKW+0CTZDQ1uBv5whptCnc8e8HeRRJxRs0kmm/Qfw== + dependencies: + depd "~1.1.2" + inherits "2.0.4" + setprototypeof "1.1.1" + statuses ">= 1.5.0 < 2" + toidentifier "1.0.0" + http-signature@~1.2.0: version "1.2.0" resolved "https://registry.yarnpkg.com/http-signature/-/http-signature-1.2.0.tgz#9aecd925114772f3d95b65a60abb8f7c18fbace1" @@ -1401,6 +1568,13 @@ https-proxy-agent@^5.0.0: agent-base "6" debug "4" +iconv-lite@0.4.24: + version "0.4.24" + resolved "https://registry.yarnpkg.com/iconv-lite/-/iconv-lite-0.4.24.tgz#2022b4b25fbddc21d2f524974a474aafe733908b" + integrity sha512-v3MXnZAcvnywkTUEZomIActle7RXXeedOR31wwl7VlyoXO4Qi9arvSenNQWne1TcRwhCL1HwLI21bEqdpj8/rA== + dependencies: + safer-buffer ">= 2.1.2 < 3" + ieee754@1.1.13, ieee754@^1.1.4: version "1.1.13" resolved "https://registry.yarnpkg.com/ieee754/-/ieee754-1.1.13.tgz#ec168558e95aa181fd87d37f55c32bbcb6708b84" @@ -1431,7 +1605,7 @@ inherits@2.0.3: resolved "https://registry.yarnpkg.com/inherits/-/inherits-2.0.3.tgz#633c2c83e3da42a502f52466022480f4208261de" integrity sha1-Yzwsg+PaQqUC9SRmAiSA9CCCYd4= -inherits@^2.0.1, inherits@^2.0.3, inherits@~2.0.1, inherits@~2.0.3: +inherits@2.0.4, inherits@^2.0.1, inherits@^2.0.3, inherits@~2.0.1, inherits@~2.0.3: version "2.0.4" resolved "https://registry.yarnpkg.com/inherits/-/inherits-2.0.4.tgz#0fa2c64f932917c3433a0ded55363aae37416b7c" integrity sha512-k/vGaX4/Yla3WzyMCvTQOXYeIHvqOKtnqBduzTHpzpQZzAskKMhZ2K+EnBiSM9zGSoIFeMpXKxa4dYeZIQqewQ== @@ -1449,6 +1623,11 @@ into-stream@^5.1.1: from2 "^2.3.0" p-is-promise "^3.0.0" +ipaddr.js@1.9.1: + version "1.9.1" + resolved "https://registry.yarnpkg.com/ipaddr.js/-/ipaddr.js-1.9.1.tgz#bff38543eeb8984825079ff3a2a8e6cbd46781b3" + integrity sha512-0KI/607xoxSToH7GjN1FfSbLoU0+btTicjsQSWQlh/hZykN8KpmMf7uYwPW3R+akZ6R/w18ZlXSHBYXiYUPO3g== + is-arrayish@^0.3.1: version "0.3.2" resolved "https://registry.yarnpkg.com/is-arrayish/-/is-arrayish-0.3.2.tgz#4574a2ae56f7ab206896fb431eaeed066fdf8f03" @@ -1764,11 +1943,26 @@ make-dir@^3.0.0: dependencies: semver "^6.0.0" +media-typer@0.3.0: + version "0.3.0" + resolved "https://registry.yarnpkg.com/media-typer/-/media-typer-0.3.0.tgz#8710d7af0aa626f8fffa1ce00168545263255748" + integrity sha1-hxDXrwqmJvj/+hzgAWhUUmMlV0g= + +merge-descriptors@1.0.1: + version "1.0.1" + resolved "https://registry.yarnpkg.com/merge-descriptors/-/merge-descriptors-1.0.1.tgz#b00aaa556dd8b44568150ec9d1b953f3f90cbb61" + integrity sha1-sAqqVW3YtEVoFQ7J0blT8/kMu2E= + merge2@^1.3.0: version "1.4.1" resolved "https://registry.yarnpkg.com/merge2/-/merge2-1.4.1.tgz#4368892f885e907455a6fd7dc55c0c9d404990ae" integrity sha512-8q7VEgMJW4J8tcfVPy8g09NcQwZdbwFEqhe/WZkoIzjn/3TGDwtOCYtXGxA3O8tPzpczCCDgv+P2P5y00ZJOOg== +methods@~1.1.2: + version "1.1.2" + resolved "https://registry.yarnpkg.com/methods/-/methods-1.1.2.tgz#5529a4d67654134edcc5266656835b0f851afcee" + integrity sha1-VSmk1nZUE07cxSZmVoNbD4Ua/O4= + micromatch@^4.0.2: version "4.0.2" resolved "https://registry.yarnpkg.com/micromatch/-/micromatch-4.0.2.tgz#4fcb0999bf9fbc2fcbdd212f6d629b9a56c39259" @@ -1782,6 +1976,11 @@ mime-db@1.44.0: resolved "https://registry.yarnpkg.com/mime-db/-/mime-db-1.44.0.tgz#fa11c5eb0aca1334b4233cb4d52f10c5a6272f92" integrity sha512-/NOTfLrsPBVeH7YtFPgsVWveuL+4SjjYxaQ1xtM1KMFj7HdxlBlxeyNLzhyJVx7r4rZGJAZ/6lkKCitSc/Nmpg== +mime-db@1.48.0: + version "1.48.0" + resolved "https://registry.yarnpkg.com/mime-db/-/mime-db-1.48.0.tgz#e35b31045dd7eada3aaad537ed88a33afbef2d1d" + integrity sha512-FM3QwxV+TnZYQ2aRqhlKBMHxk10lTbMt3bBkMAp54ddrNeVSfcQYOOKuGuy3Ddrm38I04If834fOUSq1yzslJQ== + mime-types@^2.1.12, mime-types@~2.1.19: version "2.1.27" resolved "https://registry.yarnpkg.com/mime-types/-/mime-types-2.1.27.tgz#47949f98e279ea53119f5722e0f34e529bec009f" @@ -1789,6 +1988,18 @@ mime-types@^2.1.12, mime-types@~2.1.19: dependencies: mime-db "1.44.0" +mime-types@~2.1.24: + version "2.1.31" + resolved "https://registry.yarnpkg.com/mime-types/-/mime-types-2.1.31.tgz#a00d76b74317c61f9c2db2218b8e9f8e9c5c9e6b" + integrity sha512-XGZnNzm3QvgKxa8dpzyhFTHmpP3l5YNusmne07VUOXxou9CqUqYa/HBy124RqtVh/O2pECas/MOcsDgpilPOPg== + dependencies: + mime-db "1.48.0" + +mime@1.6.0: + version "1.6.0" + resolved "https://registry.yarnpkg.com/mime/-/mime-1.6.0.tgz#32cd9e5c64553bd58d19a568af452acff04981b1" + integrity sha512-x0Vn8spI+wuJ1O6S7gnbaQg8Pxh4NNHb7KSINmEWKiPE4RKOplvijn+NkmYmmRgP68mc70j2EbeTFRsrswaQeg== + mime@^2.2.0: version "2.4.6" resolved "https://registry.yarnpkg.com/mime/-/mime-2.4.6.tgz#e5b407c90db442f2beb5b162373d07b69affa4d1" @@ -1833,6 +2044,11 @@ ms@2.0.0: resolved "https://registry.yarnpkg.com/ms/-/ms-2.0.0.tgz#5608aeadfc00be6c2901df5f9861788de0d597c8" integrity sha1-VgiurfwAvmwpAd9fmGF4jeDVl8g= +ms@2.1.1: + version "2.1.1" + resolved "https://registry.yarnpkg.com/ms/-/ms-2.1.1.tgz#30a5864eb3ebb0a66f2ebe6d727af06a09d86e0a" + integrity sha512-tgp+dl5cGk28utYktBsrFqA7HKgrhgPsg6Z/EfhWI4gl1Hwq8B/GmY/0oXZ6nF8hDVesS/FpnYaD/kOWhYQvyg== + ms@^2.1.1: version "2.1.2" resolved "https://registry.yarnpkg.com/ms/-/ms-2.1.2.tgz#d09d1f357b443f493382a8eb3ccd183872ae6009" @@ -1846,6 +2062,11 @@ multistream@^2.1.1: inherits "^2.0.1" readable-stream "^2.0.5" +negotiator@0.6.2: + version "0.6.2" + resolved "https://registry.yarnpkg.com/negotiator/-/negotiator-0.6.2.tgz#feacf7ccf525a77ae9634436a64883ffeca346fb" + integrity sha512-hZXc7K2e+PgeI1eDBe/10Ard4ekbfrrqG8Ep+8Jmf4JID2bNg7NvCPOZN+kfF574pFQI7mum2AUqDidoKqcTOw== + node-fetch@^2.3.0, node-fetch@^2.6.0: version "2.6.0" resolved "https://registry.yarnpkg.com/node-fetch/-/node-fetch-2.6.0.tgz#e633456386d4aa55863f676a7ab0daa8fdecb0fd" @@ -1899,6 +2120,13 @@ object-hash@^2.0.1: resolved "https://registry.yarnpkg.com/object-hash/-/object-hash-2.0.3.tgz#d12db044e03cd2ca3d77c0570d87225b02e1e6ea" integrity sha512-JPKn0GMu+Fa3zt3Bmr66JhokJU5BaNBIh4ZeTlaCBzrBsOeXzwcKKAK1tbLiPKgvwmPXsDvvLHoWh5Bm7ofIYg== +on-finished@~2.3.0: + version "2.3.0" + resolved "https://registry.yarnpkg.com/on-finished/-/on-finished-2.3.0.tgz#20f1336481b083cd75337992a16971aa2d906947" + integrity sha1-IPEzZIGwg811M3mSoWlxqi2QaUc= + dependencies: + ee-first "1.1.1" + once@^1.3.1, once@^1.4.0: version "1.4.0" resolved "https://registry.yarnpkg.com/once/-/once-1.4.0.tgz#583b1aa775961d4b113ac17d9c50baef9dd76bd1" @@ -1974,6 +2202,11 @@ package-json@^6.3.0: registry-url "^5.0.0" semver "^6.2.0" +parseurl@~1.3.3: + version "1.3.3" + resolved "https://registry.yarnpkg.com/parseurl/-/parseurl-1.3.3.tgz#9da19e7bee8d12dff0513ed5b76957793bc2e8d4" + integrity sha512-CiyeOxFT/JZyN5m0z9PfXw4SCBJ6Sygz1Dpl0wqjlhDEGGBP1GnsUVEL0p63hoG1fcj3fHynXi9NYO4nWOL+qQ== + path-exists@^4.0.0: version "4.0.0" resolved "https://registry.yarnpkg.com/path-exists/-/path-exists-4.0.0.tgz#513bdbe2d3b95d7762e8c1137efa195c6c61b5b3" @@ -1984,6 +2217,11 @@ path-parse@^1.0.6: resolved "https://registry.yarnpkg.com/path-parse/-/path-parse-1.0.6.tgz#d62dbb5679405d72c4737ec58600e9ddcf06d24c" integrity sha512-GSmOT2EbHrINBf9SR7CDELwlJ8AENk3Qn7OikK4nFYAu3Ote2+JYNVvkpAEQm3/TLNEJFD/xZJjzyxg3KBWOzw== +path-to-regexp@0.1.7: + version "0.1.7" + resolved "https://registry.yarnpkg.com/path-to-regexp/-/path-to-regexp-0.1.7.tgz#df604178005f522f15eb4490e7247a1bfaa67f8c" + integrity sha1-32BBeABfUi8V60SQ5yR6G/qmf4w= + path-type@^4.0.0: version "4.0.0" resolved "https://registry.yarnpkg.com/path-type/-/path-type-4.0.0.tgz#84ed01c0a7ba380afe09d90a8c180dcd9d03043b" @@ -2079,6 +2317,14 @@ protobufjs@^6.8.6, protobufjs@^6.9.0: "@types/node" "^13.7.0" long "^4.0.0" +proxy-addr@~2.0.5: + version "2.0.7" + resolved "https://registry.yarnpkg.com/proxy-addr/-/proxy-addr-2.0.7.tgz#f19fe69ceab311eeb94b42e70e8c2070f9ba1025" + integrity sha512-llQsMLSUDUPT44jdrU/O37qlnifitDP+ZwrmmZcoSKyLKvtZxpyV0n2/bD/N4tBAAZ/gJEdZU7KMraoK1+XYAg== + dependencies: + forwarded "0.2.0" + ipaddr.js "1.9.1" + psl@^1.1.28, psl@^1.1.33: version "1.8.0" resolved "https://registry.yarnpkg.com/psl/-/psl-1.8.0.tgz#9326f8bcfb013adcc005fdff056acce020e51c24" @@ -2114,6 +2360,11 @@ pupa@^2.0.1: dependencies: escape-goat "^2.0.0" +qs@6.7.0: + version "6.7.0" + resolved "https://registry.yarnpkg.com/qs/-/qs-6.7.0.tgz#41dc1a015e3d581f1621776be31afb2876a9b1bc" + integrity sha512-VCdBRNFTX1fyE7Nb6FYoURo/SPe62QCaAyzJvUjwRaIsc+NePBEniHlvxFmmX56+HZphIGtV0XeCirBtpDrTyQ== + qs@~6.5.2: version "6.5.2" resolved "https://registry.yarnpkg.com/qs/-/qs-6.5.2.tgz#cb3ae806e8740444584ef154ce8ee98d403f3e36" @@ -2129,6 +2380,21 @@ querystringify@^2.1.1: resolved "https://registry.yarnpkg.com/querystringify/-/querystringify-2.2.0.tgz#3345941b4153cb9d082d8eee4cda2016a9aef7f6" integrity sha512-FIqgj2EUvTa7R50u0rGsyTftzjYmv/a3hO345bZNrqabNqjtgiDMgmo4mkUjd+nzU5oF3dClKqFIPUKybUyqoQ== +range-parser@~1.2.1: + version "1.2.1" + resolved "https://registry.yarnpkg.com/range-parser/-/range-parser-1.2.1.tgz#3cf37023d199e1c24d1a55b84800c2f3e6468031" + integrity sha512-Hrgsx+orqoygnmhFbKaHE6c296J+HTAQXoxEF6gNupROmmGJRoyzfG3ccAveqCBrwr/2yxQ5BVd/GTl5agOwSg== + +raw-body@2.4.0: + version "2.4.0" + resolved "https://registry.yarnpkg.com/raw-body/-/raw-body-2.4.0.tgz#a1ce6fb9c9bc356ca52e89256ab59059e13d0332" + integrity sha512-4Oz8DUIwdvoa5qMJelxipzi/iJIi40O5cGV1wNYp5hvZP8ZN0T+jiNkL0QepXs+EsQ9XJ8ipEDoiH70ySUJP3Q== + dependencies: + bytes "3.1.0" + http-errors "1.7.2" + iconv-lite "0.4.24" + unpipe "1.0.0" + rc@^1.2.8: version "1.2.8" resolved "https://registry.yarnpkg.com/rc/-/rc-1.2.8.tgz#cd924bf5200a075b83c188cd6b9e211b7fc0d3ed" @@ -2292,17 +2558,17 @@ run-parallel@^1.1.9: resolved "https://registry.yarnpkg.com/run-parallel/-/run-parallel-1.1.9.tgz#c9dd3a7cf9f4b2c4b6244e173a6ed866e61dd679" integrity sha512-DEqnSRTDw/Tc3FXf49zedI638Z9onwUotBMiUFKmrO2sdFKIbXamXGQ3Axd4qgphxKB4kw/qP1w5kTxnfU1B9Q== +safe-buffer@5.1.2, safe-buffer@~5.1.0, safe-buffer@~5.1.1, safe-buffer@~5.1.2: + version "5.1.2" + resolved "https://registry.yarnpkg.com/safe-buffer/-/safe-buffer-5.1.2.tgz#991ec69d296e0313747d59bdfd2b745c35f8828d" + integrity sha512-Gd2UZBJDkXlY7GbJxfsE8/nvKkUEU1G38c1siN6QP6a9PT9MmHB8GnpscSmMJSoF8LOIrt8ud/wPtojys4G6+g== + safe-buffer@^5.0.1, safe-buffer@^5.1.2, safe-buffer@~5.2.0: version "5.2.1" resolved "https://registry.yarnpkg.com/safe-buffer/-/safe-buffer-5.2.1.tgz#1eaf9fa9bdb1fdd4ec75f58f9cdb4e6b7827eec6" integrity sha512-rp3So07KcdmmKbGvgaNxQSJr7bGVSVk5S9Eq1F+ppbRo70+YeaDxkw5Dd8NPN+GD6bjnYm2VuPuCXmpuYvmCXQ== -safe-buffer@~5.1.0, safe-buffer@~5.1.1, safe-buffer@~5.1.2: - version "5.1.2" - resolved "https://registry.yarnpkg.com/safe-buffer/-/safe-buffer-5.1.2.tgz#991ec69d296e0313747d59bdfd2b745c35f8828d" - integrity sha512-Gd2UZBJDkXlY7GbJxfsE8/nvKkUEU1G38c1siN6QP6a9PT9MmHB8GnpscSmMJSoF8LOIrt8ud/wPtojys4G6+g== - -safer-buffer@^2.0.2, safer-buffer@^2.1.0, safer-buffer@~2.1.0: +"safer-buffer@>= 2.1.2 < 3", safer-buffer@^2.0.2, safer-buffer@^2.1.0, safer-buffer@~2.1.0: version "2.1.2" resolved "https://registry.yarnpkg.com/safer-buffer/-/safer-buffer-2.1.2.tgz#44fa161b0187b9549dd84bb91802f9bd8385cd6a" integrity sha512-YZo3K82SD7Riyi0E1EQPojLz7kpepnSQI9IyPbHHg1XXXevb5dJI7tpyN2ADxGcQbHG7vcyRHk0cbwqcQriUtg== @@ -2339,11 +2605,45 @@ semver@^7.1.3: resolved "https://registry.yarnpkg.com/semver/-/semver-7.3.2.tgz#604962b052b81ed0786aae84389ffba70ffd3938" integrity sha512-OrOb32TeeambH6UrhtShmF7CRDqhL6/5XpPNp2DuRH6+9QLw/orhp72j87v8Qa1ScDkvrrBNpZcDejAirJmfXQ== +send@0.17.1: + version "0.17.1" + resolved "https://registry.yarnpkg.com/send/-/send-0.17.1.tgz#c1d8b059f7900f7466dd4938bdc44e11ddb376c8" + integrity sha512-BsVKsiGcQMFwT8UxypobUKyv7irCNRHk1T0G680vk88yf6LBByGcZJOTJCrTP2xVN6yI+XjPJcNuE3V4fT9sAg== + dependencies: + debug "2.6.9" + depd "~1.1.2" + destroy "~1.0.4" + encodeurl "~1.0.2" + escape-html "~1.0.3" + etag "~1.8.1" + fresh "0.5.2" + http-errors "~1.7.2" + mime "1.6.0" + ms "2.1.1" + on-finished "~2.3.0" + range-parser "~1.2.1" + statuses "~1.5.0" + +serve-static@1.14.1: + version "1.14.1" + resolved "https://registry.yarnpkg.com/serve-static/-/serve-static-1.14.1.tgz#666e636dc4f010f7ef29970a88a674320898b2f9" + integrity sha512-JMrvUwE54emCYWlTI+hGrGv5I8dEwmco/00EvkzIIsR7MqrHonbD9pO2MOfFnpFntl7ecpZs+3mW+XbQZu9QCg== + dependencies: + encodeurl "~1.0.2" + escape-html "~1.0.3" + parseurl "~1.3.3" + send "0.17.1" + set-blocking@^2.0.0: version "2.0.0" resolved "https://registry.yarnpkg.com/set-blocking/-/set-blocking-2.0.0.tgz#045f9782d011ae9a6803ddd382b24392b3d890f7" integrity sha1-BF+XgtARrppoA93TgrJDkrPYkPc= +setprototypeof@1.1.1: + version "1.1.1" + resolved "https://registry.yarnpkg.com/setprototypeof/-/setprototypeof-1.1.1.tgz#7e95acb24aa92f5885e0abef5ba131330d4ae683" + integrity sha512-JvdAWfbXeIGaZ9cILp38HntZSFSo3mWg6xGcJJsd+d4aRMOqauag1C63dJfDw7OaMYwEbHMOxEZ1lqVRYP2OAw== + signal-exit@^3.0.2: version "3.0.3" resolved "https://registry.yarnpkg.com/signal-exit/-/signal-exit-3.0.3.tgz#a1410c2edd8f077b08b4e253c8eacfcaf057461c" @@ -2391,6 +2691,11 @@ stack-trace@0.0.x: resolved "https://registry.yarnpkg.com/stack-trace/-/stack-trace-0.0.10.tgz#547c70b347e8d32b4e108ea1a2a159e5fdde19c0" integrity sha1-VHxws0fo0ytOEI6hoqFZ5f3eGcA= +"statuses@>= 1.5.0 < 2", statuses@~1.5.0: + version "1.5.0" + resolved "https://registry.yarnpkg.com/statuses/-/statuses-1.5.0.tgz#161c7dac177659fd9811f43771fa99381478628c" + integrity sha1-Fhx9rBd2Wf2YEfQ3cfqZOBR4Yow= + stream-browserify@^2.0.2: version "2.0.2" resolved "https://registry.yarnpkg.com/stream-browserify/-/stream-browserify-2.0.2.tgz#87521d38a44aa7ee91ce1cd2a47df0cb49dd660b" @@ -2513,6 +2818,11 @@ to-regex-range@^5.0.1: dependencies: is-number "^7.0.0" +toidentifier@1.0.0: + version "1.0.0" + resolved "https://registry.yarnpkg.com/toidentifier/-/toidentifier-1.0.0.tgz#7e1be3470f1e77948bc43d94a3c8f4d7752ba553" + integrity sha512-yaOH/Pk/VEhBWWTlhI+qXxDFXlejDGcQipMlyxda9nthulaxLZUNcUqFxokp0vcYnvteJln5FNQDRrxj3YcbVw== + touch@^3.1.0: version "3.1.0" resolved "https://registry.yarnpkg.com/touch/-/touch-3.1.0.tgz#fe365f5f75ec9ed4e56825e0bb76d24ab74af83b" @@ -2581,6 +2891,14 @@ type-fest@^0.8.1: resolved "https://registry.yarnpkg.com/type-fest/-/type-fest-0.8.1.tgz#09e249ebde851d3b1e48d27c105444667f17b83d" integrity sha512-4dbzIzqvjtgiM5rw1k5rEHtBANKmdudhGyBEajN01fEyhaAIhsoKNy6y7+IN93IfpFtwY9iqi7kD+xwKhQsNJA== +type-is@~1.6.17, type-is@~1.6.18: + version "1.6.18" + resolved "https://registry.yarnpkg.com/type-is/-/type-is-1.6.18.tgz#4e552cd05df09467dcbc4ef739de89f2cf37c131" + integrity sha512-TkRKr9sUTxEH8MdfuCSP7VizJyzRNMjj2J2do2Jr3Kym598JVdEksuzPQCnlFPW4ky9Q+iA+ma9BGm06XQBy8g== + dependencies: + media-typer "0.3.0" + mime-types "~2.1.24" + typedarray-to-buffer@^3.1.5: version "3.1.5" resolved "https://registry.yarnpkg.com/typedarray-to-buffer/-/typedarray-to-buffer-3.1.5.tgz#a97ee7a9ff42691b9f783ff1bc5112fe3fca9080" @@ -2636,6 +2954,11 @@ universalify@^1.0.0: resolved "https://registry.yarnpkg.com/universalify/-/universalify-1.0.0.tgz#b61a1da173e8435b2fe3c67d29b9adf8594bd16d" integrity sha512-rb6X1W158d7pRQBg5gkR8uPaSfiids68LTJQYOtEUhoJUWBdaQHsuT/EUduxXYxcrt4r5PJ4fuHW1MHT6p0qug== +unpipe@1.0.0, unpipe@~1.0.0: + version "1.0.0" + resolved "https://registry.yarnpkg.com/unpipe/-/unpipe-1.0.0.tgz#b2bf4ee8514aae6165b4817829d21b2ef49904ec" + integrity sha1-sr9O6FFKrmFltIF4KdIbLvSZBOw= + update-notifier@^4.0.0: version "4.1.1" resolved "https://registry.yarnpkg.com/update-notifier/-/update-notifier-4.1.1.tgz#895fc8562bbe666179500f9f2cebac4f26323746" @@ -2705,6 +3028,11 @@ util@^0.11.1: dependencies: inherits "2.0.3" +utils-merge@1.0.1: + version "1.0.1" + resolved "https://registry.yarnpkg.com/utils-merge/-/utils-merge-1.0.1.tgz#9f95710f50a267947b2ccc124741c1028427e713" + integrity sha1-n5VxD1CiZ5R7LMwSR0HBAoQn5xM= + uuid-parse@^1.1.0: version "1.1.0" resolved "https://registry.yarnpkg.com/uuid-parse/-/uuid-parse-1.1.0.tgz#7061c5a1384ae0e1f943c538094597e1b5f3a65b" @@ -2735,6 +3063,11 @@ validator@^9.4.1: resolved "https://registry.yarnpkg.com/validator/-/validator-9.4.1.tgz#abf466d398b561cd243050112c6ff1de6cc12663" integrity sha512-YV5KjzvRmSyJ1ee/Dm5UED0G+1L4GZnLN3w6/T+zZm8scVua4sOhYKWTUrKa0H/tMiJyO9QLHMPN+9mB/aMunA== +vary@~1.1.2: + version "1.1.2" + resolved "https://registry.yarnpkg.com/vary/-/vary-1.1.2.tgz#2299f02c6ded30d4a5961b0b9f74524a18f634fc" + integrity sha1-IpnwLG3tMNSllhsLn3RSShj2NPw= + verror@1.10.0: version "1.10.0" resolved "https://registry.yarnpkg.com/verror/-/verror-1.10.0.tgz#3a105ca17053af55d6e270c1f8288682e18da400" From d117cec09e25c52ea69807a4d4db051505a3b934 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 07:42:21 +0300 Subject: [PATCH 21/21] js-executor: partitions_consumed_concurrently marked as experimental --- msa/js-executor/config/custom-environment-variables.yml | 2 +- msa/js-executor/config/default.yml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 4185e65835..fdf261e25c 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -29,7 +29,7 @@ kafka: acks: "TB_KAFKA_ACKS" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge batch_size: "TB_KAFKA_BATCH_SIZE" # for producer linger_ms: "TB_KAFKA_LINGER_MS" # for producer - partitions_consumed_concurrently: "TB_KAFKA_PARTITIONS_CONSUMED_CONCURRENTLY" # increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency + partitions_consumed_concurrently: "TB_KAFKA_PARTITIONS_CONSUMED_CONCURRENTLY" # (EXPERIMENTAL) increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency requestTimeout: "TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS" compression: "TB_QUEUE_KAFKA_COMPRESSION" # gzip or uncompressed topic_properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index 190cf097ed..7a6e7c2469 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -29,7 +29,7 @@ kafka: acks: "1" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge batch_size: "128" # for producer linger_ms: "1" # for producer - partitions_consumed_concurrently: "1" # increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency + partitions_consumed_concurrently: "1" # (EXPERIMENTAL) increase this value if you are planning to handle more than one partition (scale up, scale down) - this will decrease the latency requestTimeout: "30000" # The default value in kafkajs is: 30000 compression: "gzip" # gzip or uncompressed topic_properties: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600;partitions:100;min.insync.replicas:1"