diff --git a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.test.ts b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.test.ts index c0d6809514..5824bf4d41 100644 --- a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.test.ts +++ b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.test.ts @@ -68,6 +68,16 @@ describe('KafkaConsumerClient', () => { expect(client).toBeInstanceOf(KafkaConsumerClient); }); + it('should not create an instance from config', () => { + const client = KafkaConsumerClient.fromConfig({ + config: new ConfigReader({}), + events: mockEvents, + logger: mockLogger, + }); + + expect(client).toBeUndefined(); + }); + it('should create a consumer for each topic from config', () => { KafkaConsumerClient.fromConfig({ config: mockConfig, @@ -92,7 +102,9 @@ describe('KafkaConsumerClient', () => { logger: mockLogger, }); - await client.start(); + expect(client).toBeDefined(); + + await client?.start(); expect(mockConsumer.start).toHaveBeenCalled(); }); @@ -111,7 +123,9 @@ describe('KafkaConsumerClient', () => { logger: mockLogger, }); - await client.shutdown(); + expect(client).toBeDefined(); + + await client?.shutdown(); expect(mockConsumer.shutdown).toHaveBeenCalled(); }); diff --git a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.ts b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.ts index 762a95dd41..735b0fc224 100644 --- a/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.ts +++ b/plugins/events-backend-module-kafka/src/publisher/KafkaConsumerClient.ts @@ -32,16 +32,21 @@ export class KafkaConsumerClient { private readonly kafka: Kafka; private readonly consumers: KafkaConsumingEventPublisher[]; - static fromConfig(env: { + static fromConfig(options: { config: Config; events: EventsService; logger: LoggerService; - }): KafkaConsumerClient { - return new KafkaConsumerClient( - env.logger, - env.events, - readConfig(env.config), - ); + }): KafkaConsumerClient | undefined { + const kafkaConfig = readConfig(options.config); + + if (!kafkaConfig) { + options.logger.info( + 'Kafka consumer not configured, skipping initialization', + ); + return undefined; + } + + return new KafkaConsumerClient(options.logger, options.events, kafkaConfig); } private constructor( diff --git a/plugins/events-backend-module-kafka/src/publisher/config.test.ts b/plugins/events-backend-module-kafka/src/publisher/config.test.ts index 6381880437..5b2dfb6f0c 100644 --- a/plugins/events-backend-module-kafka/src/publisher/config.test.ts +++ b/plugins/events-backend-module-kafka/src/publisher/config.test.ts @@ -18,11 +18,9 @@ import { readConfig } from './config'; describe('readConfig', () => { it('not configured', () => { - const config = new ConfigReader({}); + const publisherConfigs = readConfig(new ConfigReader({})); - expect(() => { - readConfig(config); - }).toThrow(); + expect(publisherConfigs).toBeUndefined(); }); it('only required fields configured', () => { @@ -57,21 +55,23 @@ describe('readConfig', () => { const publisherConfigs = readConfig(config); - expect(publisherConfigs.kafkaConsumerConfigs.length).toBe(2); + expect(publisherConfigs).toBeDefined(); - expect(publisherConfigs.kafkaConfig.clientId).toEqual('backstage-events'); - expect(publisherConfigs.kafkaConfig.brokers).toEqual([ + expect(publisherConfigs?.kafkaConsumerConfigs.length).toBe(2); + + expect(publisherConfigs?.kafkaConfig.clientId).toEqual('backstage-events'); + expect(publisherConfigs?.kafkaConfig.brokers).toEqual([ 'kafka1:9092', 'kafka2:9092', ]); - expect(publisherConfigs.kafkaConsumerConfigs[0].backstageTopic).toEqual( + expect(publisherConfigs?.kafkaConsumerConfigs[0].backstageTopic).toEqual( 'fake1', ); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.groupId, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.groupId, ).toEqual('my-group'); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerSubscribeTopics.topics, + publisherConfigs?.kafkaConsumerConfigs[0].consumerSubscribeTopics.topics, ).toEqual(['topic-A']); }); @@ -132,23 +132,25 @@ describe('readConfig', () => { const publisherConfigs = readConfig(config); + expect(publisherConfigs).toBeDefined(); + // Client configuration - expect(publisherConfigs.kafkaConfig.clientId).toEqual('backstage-events'); - expect(publisherConfigs.kafkaConfig.brokers).toEqual([ + expect(publisherConfigs?.kafkaConfig.clientId).toEqual('backstage-events'); + expect(publisherConfigs?.kafkaConfig.brokers).toEqual([ 'kafka1:9092', 'kafka2:9092', ]); - expect(publisherConfigs.kafkaConfig.ssl).toBeTruthy(); - expect(publisherConfigs.kafkaConfig.sasl).toStrictEqual({ + expect(publisherConfigs?.kafkaConfig.ssl).toBeTruthy(); + expect(publisherConfigs?.kafkaConfig.sasl).toStrictEqual({ mechanism: 'plain', username: 'username', password: 'password', }); - expect(publisherConfigs.kafkaConfig.authenticationTimeout).toBe(20000); - expect(publisherConfigs.kafkaConfig.connectionTimeout).toBe(1500); - expect(publisherConfigs.kafkaConfig.requestTimeout).toBe(20000); - expect(publisherConfigs.kafkaConfig.enforceRequestTimeout).toBeFalsy(); - expect(publisherConfigs.kafkaConfig.retry).toStrictEqual({ + expect(publisherConfigs?.kafkaConfig.authenticationTimeout).toBe(20000); + expect(publisherConfigs?.kafkaConfig.connectionTimeout).toBe(1500); + expect(publisherConfigs?.kafkaConfig.requestTimeout).toBe(20000); + expect(publisherConfigs?.kafkaConfig.enforceRequestTimeout).toBeFalsy(); + expect(publisherConfigs?.kafkaConfig.retry).toStrictEqual({ maxRetryTime: '20000', initialRetryTime: '200', factor: '0.4', @@ -157,41 +159,42 @@ describe('readConfig', () => { }); // Consumer configuration - expect(publisherConfigs.kafkaConsumerConfigs.length).toBe(2); - expect(publisherConfigs.kafkaConsumerConfigs[0].backstageTopic).toEqual( + expect(publisherConfigs?.kafkaConsumerConfigs.length).toBe(2); + expect(publisherConfigs?.kafkaConsumerConfigs[0].backstageTopic).toEqual( 'fake1', ); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.groupId, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.groupId, ).toEqual('my-group'); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerSubscribeTopics.topics, + publisherConfigs?.kafkaConsumerConfigs[0].consumerSubscribeTopics.topics, ).toEqual(['topic-A']); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.sessionTimeout, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.sessionTimeout, ).toBe(20000); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.rebalanceTimeout, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.rebalanceTimeout, ).toBe(50000); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.heartbeatInterval, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig + .heartbeatInterval, ).toBe(2000); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.metadataMaxAge, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.metadataMaxAge, ).toBe(400000); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig .maxBytesPerPartition, ).toBe(50000); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.minBytes, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.minBytes, ).toBe(2); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.maxBytes, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.maxBytes, ).toBe(500000); expect( - publisherConfigs.kafkaConsumerConfigs[0].consumerConfig.maxWaitTimeInMs, + publisherConfigs?.kafkaConsumerConfigs[0].consumerConfig.maxWaitTimeInMs, ).toBe(4000); }); }); diff --git a/plugins/events-backend-module-kafka/src/publisher/config.ts b/plugins/events-backend-module-kafka/src/publisher/config.ts index b872a2a8d9..aadf4a5166 100644 --- a/plugins/events-backend-module-kafka/src/publisher/config.ts +++ b/plugins/events-backend-module-kafka/src/publisher/config.ts @@ -36,8 +36,14 @@ export interface KafkaEventSourceConfig { const CONFIG_PREFIX_PUBLISHER = 'events.modules.kafka.kafkaConsumingEventPublisher'; -export const readConfig = (config: Config): KafkaEventSourceConfig => { - const kafkaConfig = config.getConfig(CONFIG_PREFIX_PUBLISHER); +export const readConfig = ( + config: Config, +): KafkaEventSourceConfig | undefined => { + const kafkaConfig = config.getOptionalConfig(CONFIG_PREFIX_PUBLISHER); + + if (!kafkaConfig) { + return undefined; + } const clientId = kafkaConfig.getString('clientId'); const brokers = kafkaConfig.getStringArray('brokers'); diff --git a/plugins/events-backend-module-kafka/src/service/eventsModuleKafkaConsumingEventPublisher.ts b/plugins/events-backend-module-kafka/src/service/eventsModuleKafkaConsumingEventPublisher.ts index a458a39662..b524454ca0 100644 --- a/plugins/events-backend-module-kafka/src/service/eventsModuleKafkaConsumingEventPublisher.ts +++ b/plugins/events-backend-module-kafka/src/service/eventsModuleKafkaConsumingEventPublisher.ts @@ -43,6 +43,10 @@ export const eventsModuleKafkaConsumingEventPublisher = createBackendModule({ logger, }); + if (!kafka) { + return; + } + await kafka.start(); lifecycle.addShutdownHook(async () => await kafka.shutdown());