Browse Source

Merge branch 'develop/3.4' into refactoring_test_19

pull/6900/head
nickAS21 4 years ago
parent
commit
c9ead2994a
  1. 31
      msa/js-executor/api/httpServer.ts
  2. 1
      msa/js-executor/docker/Dockerfile
  3. 2
      msa/js-executor/package.json
  4. 144
      msa/js-executor/queue/awsSqsTemplate.ts
  5. 186
      msa/js-executor/queue/kafkaTemplate.ts
  6. 88
      msa/js-executor/queue/pubSubTemplate.ts
  7. 3
      msa/js-executor/queue/queue.models.ts
  8. 73
      msa/js-executor/queue/rabbitmqTemplate.ts
  9. 97
      msa/js-executor/queue/serviceBusTemplate.ts
  10. 79
      msa/js-executor/server.ts
  11. 8
      msa/js-executor/yarn.lock

31
msa/js-executor/api/httpServer.ts

@ -16,12 +16,15 @@
import express from 'express'; import express from 'express';
import { _logger} from '../config/logger'; import { _logger} from '../config/logger';
import http from 'http';
import { Socket } from 'net';
export class HttpServer { export class HttpServer {
private logger = _logger('httpServer'); private logger = _logger('httpServer');
private app = express(); private app = express();
private server; private server: http.Server | null;
private connections: Socket[] = [];
constructor(httpPort: number) { constructor(httpPort: number) {
this.app.get('/livenessProbe', async (req, res) => { this.app.get('/livenessProbe', async (req, res) => {
@ -32,15 +35,31 @@ export class HttpServer {
}) })
this.server = this.app.listen(httpPort, () => { this.server = this.app.listen(httpPort, () => {
this.logger.info('Started http endpoint on port %s. Please, use /livenessProbe !', httpPort); this.logger.info('Started HTTP endpoint on port %s. Please, use /livenessProbe !', httpPort);
}).on('error', (error) => { }).on('error', (error) => {
this.logger.error(error); this.logger.error(error);
}); });
}
stop() { this.server.on('connection', connection => {
this.server.close(() => { this.connections.push(connection);
this.logger.info('Http server stop'); connection.on('close', () => this.connections = this.connections.filter(curr => curr !== connection));
}); });
} }
async stop() {
if (this.server) {
this.logger.info('Stopping HTTP Server...');
const _server = this.server;
this.server = null;
this.connections.forEach(curr => curr.end(() => curr.destroy()));
await new Promise<void>(
(resolve, reject) => {
_server.close((err) => {
this.logger.info('HTTP Server stopped.');
resolve();
});
}
);
}
}
} }

1
msa/js-executor/docker/Dockerfile

@ -29,6 +29,7 @@ COPY package/linux/conf ./conf
COPY package/linux/conf ./config COPY package/linux/conf ./config
COPY src/api ./api COPY src/api ./api
COPY src/queue ./queue COPY src/queue ./queue
COPY src/config ./config
COPY src/server.js ./ COPY src/server.js ./
RUN chmod a+x /tmp/*.sh \ RUN chmod a+x /tmp/*.sh \

2
msa/js-executor/package.json

@ -20,7 +20,7 @@
"config": "^3.3.7", "config": "^3.3.7",
"express": "^4.18.1", "express": "^4.18.1",
"js-yaml": "^4.1.0", "js-yaml": "^4.1.0",
"kafkajs": "^2.0.2", "kafkajs": "^2.1.0",
"long": "^5.2.0", "long": "^5.2.0",
"uuid-parse": "^1.1.0", "uuid-parse": "^1.1.0",
"uuid-random": "^1.3.2", "uuid-random": "^1.3.2",

144
msa/js-executor/queue/awsSqsTemplate.ts

@ -54,82 +54,76 @@ export class AwsSqsTemplate implements IQueue {
FifoQueue: 'true' FifoQueue: 'true'
}; };
name = 'AWS SQS';
constructor() { constructor() {
} }
async init() { async init() {
try { this.sqsClient = new SQSClient({
this.logger.info('Starting ThingsBoard JavaScript Executor Microservice...'); apiVersion: '2012-11-05',
credentials: {
this.sqsClient = new SQSClient({ accessKeyId: this.accessKeyId,
apiVersion: '2012-11-05', secretAccessKey: this.secretAccessKey
credentials: { },
accessKeyId: this.accessKeyId, region: this.region
secretAccessKey: this.secretAccessKey });
},
region: this.region
});
const queues = await this.getQueues(); const queues = await this.getQueues();
if (queues.QueueUrls) { if (queues.QueueUrls) {
queues.QueueUrls.forEach(queueUrl => { queues.QueueUrls.forEach(queueUrl => {
const delimiterPosition = queueUrl.lastIndexOf('/'); const delimiterPosition = queueUrl.lastIndexOf('/');
const queueName = queueUrl.substring(delimiterPosition + 1); const queueName = queueUrl.substring(delimiterPosition + 1);
this.queueUrls.set(queueName, queueUrl); this.queueUrls.set(queueName, queueUrl);
}); });
} }
this.parseQueueProperties(); this.parseQueueProperties();
this.requestQueueURL = this.queueUrls.get(AwsSqsTemplate.topicToSqsQueueName(this.requestTopic)) || ''; this.requestQueueURL = this.queueUrls.get(AwsSqsTemplate.topicToSqsQueueName(this.requestTopic)) || '';
if (!this.requestQueueURL) { if (!this.requestQueueURL) {
this.requestQueueURL = await this.createQueue(this.requestTopic); this.requestQueueURL = await this.createQueue(this.requestTopic);
} }
const messageProcessor = new JsInvokeMessageProcessor(this);
const messageProcessor = new JsInvokeMessageProcessor(this); const params: ReceiveMessageRequest = {
MaxNumberOfMessages: 10,
const params: ReceiveMessageRequest = { QueueUrl: this.requestQueueURL,
MaxNumberOfMessages: 10, WaitTimeSeconds: this.pollInterval / 1000
QueueUrl: this.requestQueueURL, };
WaitTimeSeconds: this.pollInterval / 1000 while (!this.stopped) {
}; let pollStartTs = new Date().getTime();
while (!this.stopped) { const messagesResponse: ReceiveMessageResult = await this.sqsClient.send(new ReceiveMessageCommand(params));
let pollStartTs = new Date().getTime(); const messages = messagesResponse.Messages;
const messagesResponse: ReceiveMessageResult = await this.sqsClient.send(new ReceiveMessageCommand(params));
const messages = messagesResponse.Messages; if (messages && messages.length > 0) {
const entries: DeleteMessageBatchRequestEntry[] = [];
if (messages && messages.length > 0) {
const entries: DeleteMessageBatchRequestEntry[] = []; messages.forEach(message => {
entries.push({
messages.forEach(message => { Id: message.MessageId,
entries.push({ ReceiptHandle: message.ReceiptHandle
Id: message.MessageId,
ReceiptHandle: message.ReceiptHandle
});
messageProcessor.onJsInvokeMessage(JSON.parse(message.Body || ''));
}); });
messageProcessor.onJsInvokeMessage(JSON.parse(message.Body || ''));
});
const deleteBatch: DeleteMessageBatchRequest = { const deleteBatch: DeleteMessageBatchRequest = {
QueueUrl: this.requestQueueURL, QueueUrl: this.requestQueueURL,
Entries: entries Entries: entries
}; };
try { try {
await this.sqsClient.send(new DeleteMessageBatchCommand(deleteBatch)) await this.sqsClient.send(new DeleteMessageBatchCommand(deleteBatch))
} catch (err: any) { } catch (err: any) {
this.logger.error("Failed to delete messages from queue.", err.message); this.logger.error("Failed to delete messages from queue.", err.message);
} }
} else { } else {
let pollDuration = new Date().getTime() - pollStartTs; let pollDuration = new Date().getTime() - pollStartTs;
if (pollDuration < this.pollInterval) { if (pollDuration < this.pollInterval) {
await sleep(this.pollInterval - pollDuration); await sleep(this.pollInterval - pollDuration);
}
} }
} }
} catch (e: any) {
this.logger.error('Failed to start ThingsBoard JavaScript Executor Microservice: %s', e.message);
this.logger.error(e.stack);
await this.exit(-1);
} }
} }
@ -187,29 +181,21 @@ export class AwsSqsTemplate implements IQueue {
return result.QueueUrl || ''; return result.QueueUrl || '';
} }
static async build(): Promise<AwsSqsTemplate> { async destroy(): Promise<void> {
const queue = new AwsSqsTemplate();
await queue.init();
return queue;
}
async exit(status: number) {
this.stopped = true; this.stopped = true;
this.logger.info('Exiting with status: %d ...', status); this.logger.info('Stopping AWS SQS resources...');
if (this.sqsClient) { if (this.sqsClient) {
this.logger.info('Stopping Aws Sqs client.') this.logger.info('Stopping AWS SQS client...');
try { try {
this.sqsClient.destroy(); const _sqsClient = this.sqsClient;
// @ts-ignore // @ts-ignore
delete this.sqsClient; delete this.sqsClient;
this.logger.info('Aws Sqs client stopped.') _sqsClient.destroy();
process.exit(status); this.logger.info('AWS SQS client stopped.');
} catch (e: any) { } catch (e: any) {
this.logger.info('Aws Sqs client stop error.'); this.logger.info('AWS SQS client stop error.');
process.exit(status);
} }
} else {
process.exit(status);
} }
this.logger.info('AWS SQS resources stopped.')
} }
} }

186
msa/js-executor/queue/kafkaTemplate.ts

@ -51,111 +51,103 @@ export class KafkaTemplate implements IQueue {
private batchMessages: TopicMessages[] = []; private batchMessages: TopicMessages[] = [];
private sendLoopInstance: NodeJS.Timeout; private sendLoopInstance: NodeJS.Timeout;
name = 'Kafka';
constructor() { constructor() {
} }
async init(): Promise<void> { async init(): Promise<void> {
try { const kafkaBootstrapServers: string = config.get('kafka.bootstrap.servers');
this.logger.info('Starting ThingsBoard JavaScript Executor Microservice...'); const requestTopic: string = config.get('request_topic');
const useConfluent = config.get('kafka.use_confluent_cloud');
const kafkaBootstrapServers: string = config.get('kafka.bootstrap.servers');
const requestTopic: string = config.get('request_topic');
const useConfluent = config.get('kafka.use_confluent_cloud');
this.logger.info('Kafka Bootstrap Servers: %s', kafkaBootstrapServers); this.logger.info('Kafka Bootstrap Servers: %s', kafkaBootstrapServers);
this.logger.info('Kafka Requests Topic: %s', requestTopic); this.logger.info('Kafka Requests Topic: %s', requestTopic);
let kafkaConfig: KafkaConfig = { let kafkaConfig: KafkaConfig = {
brokers: kafkaBootstrapServers.split(','), brokers: kafkaBootstrapServers.split(','),
logLevel: logLevel.INFO, logLevel: logLevel.INFO,
logCreator: KafkaJsWinstonLogCreator logCreator: KafkaJsWinstonLogCreator
}; };
if (this.kafkaClientId) { if (this.kafkaClientId) {
kafkaConfig['clientId'] = this.kafkaClientId; kafkaConfig['clientId'] = this.kafkaClientId;
} else { } else {
this.logger.warn('KAFKA_CLIENT_ID is undefined. Consider to define the env variable KAFKA_CLIENT_ID'); this.logger.warn('KAFKA_CLIENT_ID is undefined. Consider to define the env variable KAFKA_CLIENT_ID');
} }
kafkaConfig['requestTimeout'] = this.requestTimeout; kafkaConfig['requestTimeout'] = this.requestTimeout;
if (useConfluent) { if (useConfluent) {
kafkaConfig['sasl'] = { kafkaConfig['sasl'] = {
mechanism: config.get('kafka.confluent.sasl.mechanism') as any, mechanism: config.get('kafka.confluent.sasl.mechanism') as any,
username: config.get('kafka.confluent.username'), username: config.get('kafka.confluent.username'),
password: config.get('kafka.confluent.password') password: config.get('kafka.confluent.password')
}; };
kafkaConfig['ssl'] = true; kafkaConfig['ssl'] = true;
} }
this.parseTopicProperties(); this.parseTopicProperties();
this.kafkaClient = new Kafka(kafkaConfig); this.kafkaClient = new Kafka(kafkaConfig);
this.kafkaAdmin = this.kafkaClient.admin(); this.kafkaAdmin = this.kafkaClient.admin();
await this.kafkaAdmin.connect(); await this.kafkaAdmin.connect();
let partitions = 1; let partitions = 1;
for (let i = 0; i < this.configEntries.length; i++) { for (let i = 0; i < this.configEntries.length; i++) {
let param = this.configEntries[i]; let param = this.configEntries[i];
if (param.name === 'partitions') { if (param.name === 'partitions') {
partitions = param.value; partitions = param.value;
this.configEntries.splice(i, 1); this.configEntries.splice(i, 1);
break; break;
}
} }
}
let topics = await this.kafkaAdmin.listTopics(); let topics = await this.kafkaAdmin.listTopics();
if (!topics.includes(requestTopic)) { if (!topics.includes(requestTopic)) {
let createRequestTopicResult = await this.createTopic(requestTopic, partitions); let createRequestTopicResult = await this.createTopic(requestTopic, partitions);
if (createRequestTopicResult) { if (createRequestTopicResult) {
this.logger.info('Created new topic: %s', requestTopic); this.logger.info('Created new topic: %s', requestTopic);
}
} }
}
this.consumer = this.kafkaClient.consumer({groupId: 'js-executor-group'}); this.consumer = this.kafkaClient.consumer({groupId: 'js-executor-group'});
this.producer = this.kafkaClient.producer({createPartitioner: Partitioners.DefaultPartitioner}); this.producer = this.kafkaClient.producer({createPartitioner: Partitioners.DefaultPartitioner});
const {CRASH} = this.consumer.events;
this.consumer.on(CRASH, e => { const {CRASH} = this.consumer.events;
this.logger.error(`Got consumer CRASH event, should restart: ${e.payload.restart}`);
if (!e.payload.restart) {
this.logger.error('Going to exit due to not retryable error!');
this.exit(-1);
}
});
const messageProcessor = new JsInvokeMessageProcessor(this); this.consumer.on(CRASH, async (e) => {
await this.consumer.connect(); this.logger.error(`Got consumer CRASH event, should restart: ${e.payload.restart}`);
await this.producer.connect(); if (!e.payload.restart) {
this.sendLoopWithLinger(); this.logger.error('Going to exit due to not retryable error!');
await this.consumer.subscribe({topic: requestTopic}); await this.destroy();
}
this.logger.info('Started ThingsBoard JavaScript Executor Microservice.'); });
await this.consumer.run({
partitionsConsumedConcurrently: this.partitionsConsumedConcurrently,
eachMessage: async ({topic, partition, message}) => {
let headers = message.headers;
let key = message.key || new Buffer([]);
let msg = {
key: key.toString('utf8'),
data: message.value,
headers: {
data: headers
}
};
messageProcessor.onJsInvokeMessage(msg);
},
});
} catch (e: any) { const messageProcessor = new JsInvokeMessageProcessor(this);
this.logger.error('Failed to start ThingsBoard JavaScript Executor Microservice: %s', e.message); await this.consumer.connect();
this.logger.error(e.stack); await this.producer.connect();
await this.exit(-1); this.sendLoopWithLinger();
} await this.consumer.subscribe({topic: requestTopic});
}
await this.consumer.run({
partitionsConsumedConcurrently: this.partitionsConsumedConcurrently,
eachMessage: async ({topic, partition, message}) => {
let headers = message.headers;
let key = message.key || new Buffer([]);
let msg = {
key: key.toString('utf8'),
data: message.value,
headers: {
data: headers
}
};
messageProcessor.onJsInvokeMessage(msg);
},
});
}
async send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any> { async send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any> {
this.logger.debug('Pending queue response, scriptId: [%s]', scriptId); this.logger.debug('Pending queue response, scriptId: [%s]', scriptId);
@ -235,41 +227,33 @@ export class KafkaTemplate implements IQueue {
}, this.linger); }, this.linger);
} }
static async build(): Promise<KafkaTemplate> { async destroy(): Promise<void> {
const queue = new KafkaTemplate(); this.logger.info('Stopping Kafka resources...');
await queue.init();
return queue;
}
async exit(status: number): Promise<void> {
this.logger.info('Exiting with status: %d ...', status);
if (this.kafkaAdmin) { if (this.kafkaAdmin) {
this.logger.info('Stopping Kafka Admin...'); this.logger.info('Stopping Kafka Admin...');
await this.kafkaAdmin.disconnect(); const _kafkaAdmin = this.kafkaAdmin;
// @ts-ignore // @ts-ignore
delete this.kafkaAdmin; delete this.kafkaAdmin;
await _kafkaAdmin.disconnect();
this.logger.info('Kafka Admin stopped.'); this.logger.info('Kafka Admin stopped.');
} }
if (this.consumer) { if (this.consumer) {
this.logger.info('Stopping Kafka Consumer...'); this.logger.info('Stopping Kafka Consumer...');
try { try {
await this.consumer.disconnect(); const _consumer = this.consumer;
// @ts-ignore // @ts-ignore
delete this.consumer; delete this.consumer;
await _consumer.disconnect();
this.logger.info('Kafka Consumer stopped.'); this.logger.info('Kafka Consumer stopped.');
await this.disconnectProducer(); await this.disconnectProducer();
process.exit(status);
} catch (e: any) { } catch (e: any) {
this.logger.info('Kafka Consumer stop error.'); this.logger.info('Kafka Consumer stop error.');
await this.disconnectProducer(); await this.disconnectProducer();
process.exit(status);
} }
} else {
process.exit(status);
} }
this.logger.info('Kafka resources stopped.');
} }
private async disconnectProducer(): Promise<void> { private async disconnectProducer(): Promise<void> {
@ -279,13 +263,15 @@ export class KafkaTemplate implements IQueue {
this.logger.info('Stopping loop...'); this.logger.info('Stopping loop...');
clearTimeout(this.sendLoopInstance); clearTimeout(this.sendLoopInstance);
await this.sendMessagesAsBatch(); await this.sendMessagesAsBatch();
await this.producer.disconnect(); const _producer = this.producer;
// @ts-ignore // @ts-ignore
delete this.producer; delete this.producer;
await _producer.disconnect();
this.logger.info('Kafka Producer stopped.'); this.logger.info('Kafka Producer stopped.');
} catch (e) { } catch (e) {
this.logger.info('Kafka Producer stop error.'); this.logger.info('Kafka Producer stop error.');
} }
} }
} }
} }

88
msa/js-executor/queue/pubSubTemplate.ts

@ -34,56 +34,50 @@ export class PubSubTemplate implements IQueue {
private topics: string[] = []; private topics: string[] = [];
private subscriptions: string[] = []; private subscriptions: string[] = [];
name = 'Pub/Sub';
constructor() { constructor() {
} }
async init() { async init() {
try { this.pubSubClient = new PubSub({
this.logger.info('Starting ThingsBoard JavaScript Executor Microservice...'); projectId: this.projectId,
this.pubSubClient = new PubSub({ credentials: this.credentials
projectId: this.projectId, });
credentials: this.credentials
});
this.parseQueueProperties();
const topicList = await this.pubSubClient.getTopics(); this.parseQueueProperties();
if (topicList) { const topicList = await this.pubSubClient.getTopics();
topicList[0].forEach(topic => {
this.topics.push(PubSubTemplate.getName(topic.name));
});
}
const subscriptionList = await this.pubSubClient.getSubscriptions(); if (topicList) {
topicList[0].forEach(topic => {
this.topics.push(PubSubTemplate.getName(topic.name));
});
}
if (subscriptionList) { const subscriptionList = await this.pubSubClient.getSubscriptions();
topicList[0].forEach(sub => {
this.subscriptions.push(PubSubTemplate.getName(sub.name));
});
}
if (!(this.subscriptions.includes(this.requestTopic) && this.topics.includes(this.requestTopic))) { if (subscriptionList) {
await this.createTopic(this.requestTopic); topicList[0].forEach(sub => {
await this.createSubscription(this.requestTopic); this.subscriptions.push(PubSubTemplate.getName(sub.name));
} });
}
const subscription = this.pubSubClient.subscription(this.requestTopic); if (!(this.subscriptions.includes(this.requestTopic) && this.topics.includes(this.requestTopic))) {
await this.createTopic(this.requestTopic);
await this.createSubscription(this.requestTopic);
}
const messageProcessor = new JsInvokeMessageProcessor(this); const subscription = this.pubSubClient.subscription(this.requestTopic);
const messageHandler = (message: Message) => { const messageProcessor = new JsInvokeMessageProcessor(this);
messageProcessor.onJsInvokeMessage(JSON.parse(message.data.toString('utf8')));
message.ack();
};
subscription.on('message', messageHandler); const messageHandler = (message: Message) => {
messageProcessor.onJsInvokeMessage(JSON.parse(message.data.toString('utf8')));
message.ack();
};
} catch (e: any) { subscription.on('message', messageHandler);
this.logger.error('Failed to start ThingsBoard JavaScript Executor Microservice: %s', e.message);
this.logger.error(e.stack);
await this.exit(-1);
}
} }
async send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any> { async send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any> {
@ -147,29 +141,21 @@ export class PubSubTemplate implements IQueue {
} }
} }
static async build(): Promise<PubSubTemplate> { async destroy(): Promise<void> {
const queue = new PubSubTemplate(); this.logger.info('Stopping Pub/Sub resources...');
await queue.init();
return queue;
}
async exit(status: number): Promise<void> {
this.logger.info('Exiting with status: %d ...', status);
if (this.pubSubClient) { if (this.pubSubClient) {
this.logger.info('Stopping Pub/Sub client.') this.logger.info('Stopping Pub/Sub client...');
try { try {
await this.pubSubClient.close(); const _pubSubClient = this.pubSubClient;
// @ts-ignore // @ts-ignore
delete this.pubSubClient; delete this.pubSubClient;
this.logger.info('Pub/Sub client stopped.') await _pubSubClient.close();
process.exit(status); this.logger.info('Pub/Sub client stopped.');
} catch (e) { } catch (e) {
this.logger.info('Pub/Sub client stop error.'); this.logger.info('Pub/Sub client stop error.');
process.exit(status);
} }
} else {
process.exit(status);
} }
this.logger.info('Pub/Sub resources stopped.');
} }
} }

3
msa/js-executor/queue/queue.models.ts

@ -15,7 +15,8 @@
/// ///
export interface IQueue { export interface IQueue {
name: string;
init(): Promise<void>; init(): Promise<void>;
send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any>; send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any>;
exit(status: number): Promise<void>; destroy(): Promise<void>;
} }

73
msa/js-executor/queue/rabbitmqTemplate.ts

@ -44,41 +44,35 @@ export class RabbitMqTemplate implements IQueue {
private stopped = false; private stopped = false;
private topics: string[] = []; private topics: string[] = [];
name = 'RabbitMQ';
constructor() { constructor() {
} }
async init(): Promise<void> { async init(): Promise<void> {
try { const url = `amqp://${this.username}:${this.password}@${this.host}:${this.port}${this.vhost}`;
this.logger.info('Starting ThingsBoard JavaScript Executor Microservice...'); this.connection = await amqp.connect(url);
this.channel = await this.connection.createConfirmChannel();
const url = `amqp://${this.username}:${this.password}@${this.host}:${this.port}${this.vhost}`;
this.connection = await amqp.connect(url);
this.channel = await this.connection.createConfirmChannel();
this.parseQueueProperties(); this.parseQueueProperties();
await this.createQueue(this.requestTopic); await this.createQueue(this.requestTopic);
const messageProcessor = new JsInvokeMessageProcessor(this); const messageProcessor = new JsInvokeMessageProcessor(this);
while (!this.stopped) { while (!this.stopped) {
let pollStartTs = new Date().getTime(); let pollStartTs = new Date().getTime();
let message = await this.channel.get(this.requestTopic); let message = await this.channel.get(this.requestTopic);
if (message) { if (message) {
messageProcessor.onJsInvokeMessage(JSON.parse(message.content.toString('utf8'))); messageProcessor.onJsInvokeMessage(JSON.parse(message.content.toString('utf8')));
this.channel.ack(message); this.channel.ack(message);
} else { } else {
let pollDuration = new Date().getTime() - pollStartTs; let pollDuration = new Date().getTime() - pollStartTs;
if (pollDuration < this.pollInterval) { if (pollDuration < this.pollInterval) {
await sleep(this.pollInterval - pollDuration); await sleep(this.pollInterval - pollDuration);
}
} }
} }
} catch (e: any) {
this.logger.error('Failed to start ThingsBoard JavaScript Executor Microservice: %s', e.message);
this.logger.error(e.stack);
await this.exit(-1);
} }
} }
@ -114,38 +108,31 @@ export class RabbitMqTemplate implements IQueue {
return this.channel.assertQueue(topic, this.queueOptions); return this.channel.assertQueue(topic, this.queueOptions);
} }
static async build(): Promise<RabbitMqTemplate> { async destroy() {
const queue = new RabbitMqTemplate(); this.logger.info('Stopping RabbitMQ resources...');
await queue.init();
return queue;
}
async exit(status: number) {
this.logger.info('Exiting with status: %d ...', status);
if (this.channel) { if (this.channel) {
this.logger.info('Stopping RabbitMq chanel.') this.logger.info('Stopping RabbitMQ chanel...');
await this.channel.close(); const _channel = this.channel;
// @ts-ignore // @ts-ignore
delete this.channel; delete this.channel;
this.logger.info('RabbitMq chanel stopped'); await _channel.close();
this.logger.info('RabbitMQ chanel stopped');
} }
if (this.connection) { if (this.connection) {
this.logger.info('Stopping RabbitMq connection.') this.logger.info('Stopping RabbitMQ connection...')
try { try {
await this.connection.close(); const _connection = this.connection;
// @ts-ignore // @ts-ignore
delete this.connection; delete this.connection;
this.logger.info('RabbitMq client connection.') await _connection.close();
process.exit(status); this.logger.info('RabbitMQ client connection.');
} catch (e) { } catch (e) {
this.logger.info('RabbitMq connection stop error.'); this.logger.info('RabbitMQ connection stop error.');
process.exit(status);
} }
} else {
process.exit(status);
} }
this.logger.info('RabbitMQ resources stopped.')
} }
} }

97
msa/js-executor/queue/serviceBusTemplate.ts

@ -44,48 +44,42 @@ export class ServiceBusTemplate implements IQueue {
private receiver: ServiceBusReceiver; private receiver: ServiceBusReceiver;
private senderMap = new Map<string, ServiceBusSender>(); private senderMap = new Map<string, ServiceBusSender>();
name = 'Azure Service Bus';
constructor() { constructor() {
} }
async init() { async init() {
try { const connectionString = `Endpoint=sb://${this.namespaceName}.servicebus.windows.net/;SharedAccessKeyName=${this.sasKeyName};SharedAccessKey=${this.sasKey}`;
this.logger.info('Starting ThingsBoard JavaScript Executor Microservice...'); this.sbClient = new ServiceBusClient(connectionString)
this.serviceBusService = new ServiceBusAdministrationClient(connectionString);
const connectionString = `Endpoint=sb://${this.namespaceName}.servicebus.windows.net/;SharedAccessKeyName=${this.sasKeyName};SharedAccessKey=${this.sasKey}`; this.parseQueueProperties();
this.sbClient = new ServiceBusClient(connectionString)
this.serviceBusService = new ServiceBusAdministrationClient(connectionString);
this.parseQueueProperties(); const listQueues = await this.serviceBusService.listQueues();
for await (const queue of listQueues) {
this.queues.push(queue.name);
}
const listQueues = await this.serviceBusService.listQueues(); if (!this.queues.includes(this.requestTopic)) {
for await (const queue of listQueues) { await this.createQueueIfNotExist(this.requestTopic);
this.queues.push(queue.name); this.queues.push(this.requestTopic);
} }
if (!this.queues.includes(this.requestTopic)) { this.receiver = this.sbClient.createReceiver(this.requestTopic, {receiveMode: 'peekLock'});
await this.createQueueIfNotExist(this.requestTopic);
this.queues.push(this.requestTopic);
}
this.receiver = this.sbClient.createReceiver(this.requestTopic, {receiveMode: 'peekLock'}); const messageProcessor = new JsInvokeMessageProcessor(this);
const messageProcessor = new JsInvokeMessageProcessor(this); const messageHandler = async (message: ServiceBusReceivedMessage) => {
if (message) {
const messageHandler = async (message: ServiceBusReceivedMessage) => { messageProcessor.onJsInvokeMessage(message.body);
if (message) { await this.receiver.completeMessage(message);
messageProcessor.onJsInvokeMessage(message.body); }
await this.receiver.completeMessage(message); };
} const errorHandler = async (error: ProcessErrorArgs) => {
}; this.logger.error('Failed to receive message from queue.', error);
const errorHandler = async (error: ProcessErrorArgs) => { };
this.logger.error('Failed to receive message from queue.', error); this.receiver.subscribe({processMessage: messageHandler, processError: errorHandler})
};
this.receiver.subscribe({processMessage: messageHandler, processError: errorHandler})
} catch (e: any) {
this.logger.error('Failed to start ThingsBoard JavaScript Executor Microservice: %s', e.message);
this.logger.error(e.stack);
await this.exit(-1);
}
} }
async send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any> { async send(responseTopic: string, scriptId: string, rawResponse: Buffer, headers: any): Promise<any> {
@ -135,41 +129,46 @@ export class ServiceBusTemplate implements IQueue {
} }
} }
static async build(): Promise<ServiceBusTemplate> { async destroy() {
const queue = new ServiceBusTemplate();
await queue.init();
return queue;
}
async exit(status: number) {
this.logger.info('Exiting with status: %d ...', status);
this.logger.info('Stopping Azure Service Bus resources...') this.logger.info('Stopping Azure Service Bus resources...')
if (this.receiver) { if (this.receiver) {
this.logger.info('Stopping Service Bus Receiver...');
try { try {
await this.receiver.close(); const _receiver = this.receiver;
// @ts-ignore // @ts-ignore
delete this.receiver; delete this.receiver;
await _receiver.close();
this.logger.info('Service Bus Receiver stopped.');
} catch (e) { } catch (e) {
this.logger.info('Service Bus Receiver stop error.');
} }
} }
this.senderMap.forEach(k => { this.logger.info('Stopping Service Bus Senders...');
try { const senders: Promise<void>[] = [];
k.close(); this.senderMap.forEach((sender) => {
} catch (e) { senders.push(sender.close());
}
}); });
this.senderMap.clear(); this.senderMap.clear();
try {
await Promise.all(senders);
this.logger.info('Service Bus Senders stopped.');
} catch (e) {
this.logger.info('Service Bus Senders stop error.');
}
if (this.sbClient) { if (this.sbClient) {
this.logger.info('Stopping Service Bus Client...');
try { try {
await this.sbClient.close(); const _sbClient = this.sbClient;
// @ts-ignore // @ts-ignore
delete this.sbClient; delete this.sbClient;
await _sbClient.close();
this.logger.info('Service Bus Client stopped.');
} catch (e) { } catch (e) {
this.logger.info('Service Bus Client stop error.');
} }
} }
this.logger.info('Azure Service Bus resources stopped.') this.logger.info('Azure Service Bus resources stopped.')
process.exit(status);
} }
} }

79
msa/js-executor/server.ts

@ -30,57 +30,66 @@ logger.info('===CONFIG BEGIN===');
logger.info(JSON.stringify(config, null, 4)); logger.info(JSON.stringify(config, null, 4));
logger.info('===CONFIG END==='); logger.info('===CONFIG END===');
const serviceType = config.get('queue_type'); const serviceType: string = config.get('queue_type');
const httpPort = Number(config.get('http_port')); const httpPort = Number(config.get('http_port'));
let queues: IQueue; let queues: IQueue | null;
let httpServer: HttpServer; let httpServer: HttpServer | null;
(async () => { (async () => {
logger.info('Starting ThingsBoard JavaScript Executor Microservice...');
try {
queues = await createQueue(serviceType);
logger.info(`Starting ${queues.name} template...`);
await queues.init();
logger.info(`${queues.name} template started.`);
httpServer = new HttpServer(httpPort);
} catch (e: any) {
logger.error('Failed to start ThingsBoard JavaScript Executor Microservice: %s', e.message);
logger.error(e.stack);
await exit(-1);
}
})();
async function createQueue(serviceType: string): Promise<IQueue> {
switch (serviceType) { switch (serviceType) {
case 'kafka': case 'kafka':
logger.info('Starting kafka template.'); return new KafkaTemplate();
queues = await KafkaTemplate.build();
logger.info('kafka template started.');
break;
case 'pubsub': case 'pubsub':
logger.info('Starting Pub/Sub template.') return new PubSubTemplate();
queues = await PubSubTemplate.build();
logger.info('Pub/Sub template started.')
break;
case 'aws-sqs': case 'aws-sqs':
logger.info('Starting Aws Sqs template.') return new AwsSqsTemplate();
queues = await AwsSqsTemplate.build();
logger.info('Aws Sqs template started.')
break;
case 'rabbitmq': case 'rabbitmq':
logger.info('Starting RabbitMq template.') return new RabbitMqTemplate();
queues = await RabbitMqTemplate.build();
logger.info('RabbitMq template started.')
break;
case 'service-bus': case 'service-bus':
logger.info('Starting Azure Service Bus template.') return new ServiceBusTemplate();
queues = await ServiceBusTemplate.build();
logger.info('Azure Service Bus template started.')
break;
default: default:
logger.error('Unknown service type: ', serviceType); throw new Error('Unknown service type: ' + serviceType);
process.exit(-1);
} }
}
httpServer = new HttpServer(httpPort); [`SIGINT`, `SIGUSR1`, `SIGUSR2`, `uncaughtException`, `SIGTERM`].forEach((eventType) => {
})(); process.on(eventType, async () => {
logger.info(`${eventType} signal received`);
await exit(0);
})
})
process.on('SIGTERM', () => { process.on('exit', (code: number) => {
logger.info('SIGTERM signal received'); logger.info(`ThingsBoard JavaScript Executor Microservice has been stopped. Exit code: ${code}.`);
process.exit(0);
}); });
process.on('exit', async () => { async function exit(status: number) {
logger.info('Exiting with status: %d ...', status);
if (httpServer) { if (httpServer) {
httpServer.stop(); const _httpServer = httpServer;
httpServer = null;
await _httpServer.stop();
} }
if (queues) { if (queues) {
queues.exit(0); const _queues = queues;
queues = null;
await _queues.destroy();
} }
logger.info('JavaScript Executor Microservice has been stopped.'); process.exit(status);
}); }

8
msa/js-executor/yarn.lock

@ -2670,10 +2670,10 @@ jws@^4.0.0:
jwa "^2.0.0" jwa "^2.0.0"
safe-buffer "^5.0.1" safe-buffer "^5.0.1"
kafkajs@^2.0.2: kafkajs@^2.1.0:
version "2.0.2" version "2.1.0"
resolved "https://registry.yarnpkg.com/kafkajs/-/kafkajs-2.0.2.tgz#cdfc8f57aa4fd69f6d9ca1cce4ee89bbc2a3a1f9" resolved "https://registry.yarnpkg.com/kafkajs/-/kafkajs-2.1.0.tgz#32ede4e8080cc75586c5e4406eeb582fa73f7b1e"
integrity sha512-g6CM3fAenofOjR1bfOAqeZUEaSGhNtBscNokybSdW1rmIKYNwBPC9xQzwulFJm36u/xcxXUiCl/L/qfslapihA== integrity sha512-6IYiOdGWvFPbSbVB+AV3feT+A7vzw5sXm7Ze4QTfP7FRNdY8pGcpiNPvD2lfgYFD8Dm9KbMgBgTt2mf8KaIkzw==
keyv@^3.0.0: keyv@^3.0.0:
version "3.1.0" version "3.1.0"

Loading…
Cancel
Save