Browse Source
Merge pull request #9527 from dashevchenko/queue_prefix_js_executor
Added global queue prefix support for js-executor
pull/9548/head
Andrew Shvayka
3 years ago
committed by
GitHub
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
7 changed files with
12 additions and
5 deletions
-
msa/js-executor/config/custom-environment-variables.yml
-
msa/js-executor/config/default.yml
-
msa/js-executor/queue/awsSqsTemplate.ts
-
msa/js-executor/queue/kafkaTemplate.ts
-
msa/js-executor/queue/pubSubTemplate.ts
-
msa/js-executor/queue/rabbitmqTemplate.ts
-
msa/js-executor/queue/serviceBusTemplate.ts
|
|
|
@ -15,6 +15,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) |
|
|
|
queue_prefix: "TB_QUEUE_PREFIX" |
|
|
|
request_topic: "REMOTE_JS_EVAL_REQUEST_TOPIC" |
|
|
|
http_port: "HTTP_PORT" # /livenessProbe |
|
|
|
|
|
|
|
|
|
|
|
@ -16,6 +16,7 @@ |
|
|
|
|
|
|
|
queue_type: "kafka" |
|
|
|
request_topic: "js_eval.requests" |
|
|
|
queue_prefix: "" |
|
|
|
http_port: "8888" # /livenessProbe |
|
|
|
|
|
|
|
js: |
|
|
|
|
|
|
|
@ -38,7 +38,8 @@ import uuid from 'uuid-random'; |
|
|
|
export class AwsSqsTemplate implements IQueue { |
|
|
|
|
|
|
|
private logger = _logger(`awsSqsTemplate`); |
|
|
|
private requestTopic: string = config.get('request_topic'); |
|
|
|
private queuePrefix: string = config.get('queue_prefix'); |
|
|
|
private requestTopic: string = this.queuePrefix ? this.queuePrefix + "." + config.get('request_topic') : config.get('request_topic'); |
|
|
|
private accessKeyId: string = config.get('aws_sqs.access_key_id'); |
|
|
|
private secretAccessKey: string = config.get('aws_sqs.secret_access_key'); |
|
|
|
private region: string = config.get('aws_sqs.region'); |
|
|
|
|
|
|
|
@ -61,7 +61,8 @@ export class KafkaTemplate implements IQueue { |
|
|
|
|
|
|
|
async init(): Promise<void> { |
|
|
|
const kafkaBootstrapServers: string = config.get('kafka.bootstrap.servers'); |
|
|
|
const requestTopic: string = config.get('request_topic'); |
|
|
|
const queuePrefix: string = config.get('queue_prefix'); |
|
|
|
const requestTopic: string = queuePrefix ? queuePrefix + "." + config.get('request_topic') : config.get('request_topic'); |
|
|
|
const useConfluent = config.get('kafka.use_confluent_cloud'); |
|
|
|
|
|
|
|
this.logger.info('Kafka Bootstrap Servers: %s', kafkaBootstrapServers); |
|
|
|
|
|
|
|
@ -26,7 +26,8 @@ export class PubSubTemplate implements IQueue { |
|
|
|
private logger = _logger(`pubSubTemplate`); |
|
|
|
private projectId: string = config.get('pubsub.project_id'); |
|
|
|
private credentials = JSON.parse(config.get('pubsub.service_account')); |
|
|
|
private requestTopic: string = config.get('request_topic'); |
|
|
|
private queuePrefix: string = config.get('queue_prefix'); |
|
|
|
private requestTopic: string = this.queuePrefix ? this.queuePrefix + "." + config.get('request_topic') : config.get('request_topic'); |
|
|
|
private queueProperties: string = config.get('pubsub.queue_properties'); |
|
|
|
|
|
|
|
private pubSubClient: PubSub; |
|
|
|
|
|
|
|
@ -24,7 +24,8 @@ import { Options, Replies } from 'amqplib/properties'; |
|
|
|
export class RabbitMqTemplate implements IQueue { |
|
|
|
|
|
|
|
private logger = _logger(`rabbitmqTemplate`); |
|
|
|
private requestTopic: string = config.get('request_topic'); |
|
|
|
private queuePrefix: string = config.get('queue_prefix'); |
|
|
|
private requestTopic: string = this.queuePrefix ? this.queuePrefix + "." + config.get('request_topic') : config.get('request_topic'); |
|
|
|
private host = config.get('rabbitmq.host'); |
|
|
|
private port = config.get('rabbitmq.port'); |
|
|
|
private vhost = config.get('rabbitmq.virtual_host'); |
|
|
|
|
|
|
|
@ -31,7 +31,8 @@ import { |
|
|
|
export class ServiceBusTemplate implements IQueue { |
|
|
|
|
|
|
|
private logger = _logger(`serviceBusTemplate`); |
|
|
|
private requestTopic: string = config.get('request_topic'); |
|
|
|
private queuePrefix: string = config.get('queue_prefix'); |
|
|
|
private requestTopic: string = this.queuePrefix ? this.queuePrefix + "." + config.get('request_topic') : config.get('request_topic'); |
|
|
|
private namespaceName = config.get('service_bus.namespace_name'); |
|
|
|
private sasKeyName = config.get('service_bus.sas_key_name'); |
|
|
|
private sasKey = config.get('service_bus.sas_key'); |
|
|
|
|