|
|
@ -30,50 +30,57 @@ 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 | null; |
|
|
let queues: IQueue | null; |
|
|
let httpServer: HttpServer | null; |
|
|
let httpServer: HttpServer | null; |
|
|
|
|
|
|
|
|
(async () => { |
|
|
(async () => { |
|
|
logger.info('Starting ThingsBoard JavaScript Executor Microservice...'); |
|
|
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) => { |
|
|
[`SIGINT`, `SIGUSR1`, `SIGUSR2`, `uncaughtException`, `SIGTERM`].forEach((eventType) => { |
|
|
process.on(eventType, async () => { |
|
|
process.on(eventType, async () => { |
|
|
logger.info(`${eventType} signal received`); |
|
|
logger.info(`${eventType} signal received`); |
|
|
|
|
|
await exit(0); |
|
|
|
|
|
}) |
|
|
|
|
|
}) |
|
|
|
|
|
|
|
|
|
|
|
process.on('exit', (code: number) => { |
|
|
|
|
|
logger.info(`ThingsBoard JavaScript Executor Microservice has been stopped. Exit code: ${code}.`); |
|
|
|
|
|
}); |
|
|
|
|
|
|
|
|
|
|
|
async function exit(status: number) { |
|
|
|
|
|
logger.info('Exiting with status: %d ...', status); |
|
|
if (httpServer) { |
|
|
if (httpServer) { |
|
|
const _httpServer = httpServer; |
|
|
const _httpServer = httpServer; |
|
|
httpServer = null; |
|
|
httpServer = null; |
|
|
@ -82,11 +89,7 @@ let httpServer: HttpServer | null; |
|
|
if (queues) { |
|
|
if (queues) { |
|
|
const _queues = queues; |
|
|
const _queues = queues; |
|
|
queues = null; |
|
|
queues = null; |
|
|
await _queues.destroy(0); |
|
|
await _queues.destroy(); |
|
|
|
|
|
} |
|
|
|
|
|
process.exit(status); |
|
|
} |
|
|
} |
|
|
}) |
|
|
|
|
|
}) |
|
|
|
|
|
|
|
|
|
|
|
process.on('exit', (code: number) => { |
|
|
|
|
|
logger.info(`JavaScript Executor Microservice has been stopped. Exit code: ${code}.`); |
|
|
|
|
|
}); |
|
|
|
|
|
|