From 5f7f573e3312155c13d0ef0038560d9f3a516ac7 Mon Sep 17 00:00:00 2001 From: Irfan Hodzic Date: Fri, 21 Nov 2025 12:21:37 +0100 Subject: [PATCH 1/4] use for await instead of stream.on --- packages/kafka/lib/AbstractKafkaConsumer.ts | 21 ++++++--- .../test/consumer/PermissionConsumer.spec.ts | 46 +++++++++++++++++++ 2 files changed, 61 insertions(+), 6 deletions(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 0ce936b5..5e467c59 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -211,17 +211,26 @@ export abstract class AbstractKafkaConsumer< ) this.messageBatchStream.on('error', (error) => this.handlerError(error)) } else { - this.consumerStream.on('data', (message) => - this.consume( - message.topic, - message as DeserializedMessage>, - ), - ) + // biome-ignore lint/style/noNonNullAssertion: consumerStream is always created + const stream = this.consumerStream! + + this.handleSyncStream(stream).catch(this.handlerError) } this.consumerStream.on('error', (error) => this.handlerError(error)) } + private async handleSyncStream( + stream: MessagesStream, + ): Promise { + for await (const message of stream) { + await this.consume( + message.topic, + message as DeserializedMessage>, + ) + } + } + async close(): Promise { if (!this.consumerStream && !this.messageBatchStream) { // Leaving the group in case consumer joined but streams were not created diff --git a/packages/kafka/test/consumer/PermissionConsumer.spec.ts b/packages/kafka/test/consumer/PermissionConsumer.spec.ts index 5f8d9f61..ca8c27f7 100644 --- a/packages/kafka/test/consumer/PermissionConsumer.spec.ts +++ b/packages/kafka/test/consumer/PermissionConsumer.spec.ts @@ -489,4 +489,50 @@ describe('PermissionConsumer', () => { }) }) }) + + describe('sync message processing ', () => { + let publisher: PermissionPublisher + + beforeAll(() => { + publisher = new PermissionPublisher(testContext.cradle) + }) + + afterAll(async () => { + await publisher.close() + }) + + it('should process a single message', async () => { + const consumer = new PermissionConsumer(testContext.cradle) + await consumer.init() + + // When + await publisher.publish('permission-added', { id: '1', type: 'added', permissions: [] }) + + // Then + await consumer.handlerSpy.waitForMessageWithId('1', 'consumed') + expect(consumer.addedMessages).toHaveLength(1) + expect(consumer.addedMessages[0]!.value.id).toBe('1') + }) + + it('should process messages sequentially', async () => { + const consumer = new PermissionConsumer(testContext.cradle) + await consumer.init() + + // When - publish messages + await publisher.publish('permission-added', { id: '1', type: 'added', permissions: [] }) + await consumer.handlerSpy.waitForMessageWithId('1', 'consumed') + + await publisher.publish('permission-added', { id: '2', type: 'added', permissions: [] }) + await consumer.handlerSpy.waitForMessageWithId('2', 'consumed') + + await publisher.publish('permission-added', { id: '3', type: 'added', permissions: [] }) + await consumer.handlerSpy.waitForMessageWithId('3', 'consumed') + + // Verify messages were processed in order + expect(consumer.addedMessages).toHaveLength(3) + expect(consumer.addedMessages[0]!.value.id).toBe('1') + expect(consumer.addedMessages[1]!.value.id).toBe('2') + expect(consumer.addedMessages[2]!.value.id).toBe('3') + }) + }) }) From 210f209dfe3b3d086a156403ab90c1d76086dbe5 Mon Sep 17 00:00:00 2001 From: Irfan Hodzic Date: Fri, 21 Nov 2025 12:30:48 +0100 Subject: [PATCH 2/4] added a coment and consumer.close on test complete --- packages/kafka/lib/AbstractKafkaConsumer.ts | 2 ++ packages/kafka/test/consumer/PermissionConsumer.spec.ts | 6 ++++-- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 5e467c59..2f75ebd4 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -214,6 +214,8 @@ export abstract class AbstractKafkaConsumer< // biome-ignore lint/style/noNonNullAssertion: consumerStream is always created const stream = this.consumerStream! + // we are not waiting for the stream to complete + // because init() must return promised void this.handleSyncStream(stream).catch(this.handlerError) } diff --git a/packages/kafka/test/consumer/PermissionConsumer.spec.ts b/packages/kafka/test/consumer/PermissionConsumer.spec.ts index ca8c27f7..eed0cf22 100644 --- a/packages/kafka/test/consumer/PermissionConsumer.spec.ts +++ b/packages/kafka/test/consumer/PermissionConsumer.spec.ts @@ -492,6 +492,7 @@ describe('PermissionConsumer', () => { describe('sync message processing ', () => { let publisher: PermissionPublisher + let consumer: PermissionConsumer | undefined beforeAll(() => { publisher = new PermissionPublisher(testContext.cradle) @@ -499,10 +500,11 @@ describe('PermissionConsumer', () => { afterAll(async () => { await publisher.close() + await consumer?.close() }) it('should process a single message', async () => { - const consumer = new PermissionConsumer(testContext.cradle) + consumer = new PermissionConsumer(testContext.cradle) await consumer.init() // When @@ -515,7 +517,7 @@ describe('PermissionConsumer', () => { }) it('should process messages sequentially', async () => { - const consumer = new PermissionConsumer(testContext.cradle) + consumer = new PermissionConsumer(testContext.cradle) await consumer.init() // When - publish messages From 6156c5ee2219b574ab797e891009139a684cc0db Mon Sep 17 00:00:00 2001 From: Irfan Hodzic Date: Fri, 21 Nov 2025 14:39:29 +0100 Subject: [PATCH 3/4] correct tests --- .../test/consumer/PermissionConsumer.spec.ts | 248 ++++++++++++++++-- 1 file changed, 226 insertions(+), 22 deletions(-) diff --git a/packages/kafka/test/consumer/PermissionConsumer.spec.ts b/packages/kafka/test/consumer/PermissionConsumer.spec.ts index eed0cf22..85025121 100644 --- a/packages/kafka/test/consumer/PermissionConsumer.spec.ts +++ b/packages/kafka/test/consumer/PermissionConsumer.spec.ts @@ -7,6 +7,7 @@ import { KafkaHandlerConfig, type RequestContext } from '../../lib/index.ts' import { PermissionPublisher } from '../publisher/PermissionPublisher.ts' import { PERMISSION_ADDED_SCHEMA, + PERMISSION_REMOVED_SCHEMA, type PermissionAdded, TOPICS, } from '../utils/permissionSchemas.ts' @@ -490,7 +491,7 @@ describe('PermissionConsumer', () => { }) }) - describe('sync message processing ', () => { + describe('sync message processing', () => { let publisher: PermissionPublisher let consumer: PermissionConsumer | undefined @@ -498,43 +499,246 @@ describe('PermissionConsumer', () => { publisher = new PermissionPublisher(testContext.cradle) }) + beforeEach(async () => { + // Close and clear previous consumer to avoid message accumulation + if (consumer) { + await consumer.close() + consumer.clear() + } + }) + afterAll(async () => { await publisher.close() await consumer?.close() }) - it('should process a single message', async () => { - consumer = new PermissionConsumer(testContext.cradle) + it('should process messages one at a time using handleSyncStream', async () => { + // Given - track processing order and timing + const processingOrder: string[] = [] + const processingTimestamps: Record = {} + const testMessageIds = ['sync-1', 'sync-2', 'sync-3'] + + consumer = new PermissionConsumer(testContext.cradle, { + handlers: { + 'permission-added': new KafkaHandlerConfig(PERMISSION_ADDED_SCHEMA, async (message) => { + // Only track messages from this test + if (!testMessageIds.includes(message.value.id)) { + consumer!.addedMessages.push(message) + return + } + + const messageId = message.value.id + processingOrder.push(`start-${messageId}`) + processingTimestamps[messageId] = { start: Date.now(), end: 0 } + + // Simulate async work to verify sequential processing + await new Promise((resolve) => setTimeout(resolve, 50)) + + processingOrder.push(`end-${messageId}`) + processingTimestamps[messageId]!.end = Date.now() + consumer!.addedMessages.push(message) + }), + }, + }) + await consumer.init() - // When - await publisher.publish('permission-added', { id: '1', type: 'added', permissions: [] }) + // When - publish multiple messages at once + await Promise.all([ + publisher.publish('permission-added', { id: 'sync-1', type: 'added', permissions: [] }), + publisher.publish('permission-added', { id: 'sync-2', type: 'added', permissions: [] }), + publisher.publish('permission-added', { id: 'sync-3', type: 'added', permissions: [] }), + ]) + + // Then - wait for all messages to be processed + await consumer.handlerSpy.waitForMessageWithId('sync-1', 'consumed') + await consumer.handlerSpy.waitForMessageWithId('sync-2', 'consumed') + await consumer.handlerSpy.waitForMessageWithId('sync-3', 'consumed') + + // Verify messages were processed sequentially (one completes before next starts) + expect(processingOrder).toEqual([ + 'start-sync-1', + 'end-sync-1', + 'start-sync-2', + 'end-sync-2', + 'start-sync-3', + 'end-sync-3', + ]) + + // Verify each message completes before the next one starts + expect(processingTimestamps['sync-1']!.end).toBeLessThan( + processingTimestamps['sync-2']!.start, + ) + expect(processingTimestamps['sync-2']!.end).toBeLessThan( + processingTimestamps['sync-3']!.start, + ) - // Then - await consumer.handlerSpy.waitForMessageWithId('1', 'consumed') - expect(consumer.addedMessages).toHaveLength(1) - expect(consumer.addedMessages[0]!.value.id).toBe('1') + const testMessages = consumer.addedMessages.filter((m) => testMessageIds.includes(m.value.id)) + expect(testMessages).toHaveLength(3) + expect(testMessages[0]!.value.id).toBe('sync-1') + expect(testMessages[1]!.value.id).toBe('sync-2') + expect(testMessages[2]!.value.id).toBe('sync-3') }) - it('should process messages sequentially', async () => { + it('should process messages in order even when published rapidly', async () => { + // Given + const testMessageIds = ['rapid-1', 'rapid-2', 'rapid-3', 'rapid-4', 'rapid-5'] consumer = new PermissionConsumer(testContext.cradle) await consumer.init() - // When - publish messages - await publisher.publish('permission-added', { id: '1', type: 'added', permissions: [] }) - await consumer.handlerSpy.waitForMessageWithId('1', 'consumed') - - await publisher.publish('permission-added', { id: '2', type: 'added', permissions: [] }) - await consumer.handlerSpy.waitForMessageWithId('2', 'consumed') + // When - publish messages rapidly without waiting + const publishPromises = [] + for (let i = 1; i <= 5; i++) { + publishPromises.push( + publisher.publish('permission-added', { + id: `rapid-${i}`, + type: 'added', + permissions: [], + }), + ) + } + await Promise.all(publishPromises) - await publisher.publish('permission-added', { id: '3', type: 'added', permissions: [] }) - await consumer.handlerSpy.waitForMessageWithId('3', 'consumed') + // Then - wait for all messages to be processed + for (let i = 1; i <= 5; i++) { + await consumer.handlerSpy.waitForMessageWithId(`rapid-${i}`, 'consumed') + } // Verify messages were processed in order - expect(consumer.addedMessages).toHaveLength(3) - expect(consumer.addedMessages[0]!.value.id).toBe('1') - expect(consumer.addedMessages[1]!.value.id).toBe('2') - expect(consumer.addedMessages[2]!.value.id).toBe('3') + const testMessages = consumer.addedMessages.filter((m) => testMessageIds.includes(m.value.id)) + expect(testMessages).toHaveLength(5) + for (let i = 0; i < 5; i++) { + expect(testMessages[i]!.value.id).toBe(`rapid-${i + 1}`) + } + }) + + it('should ensure previous message completes before next message starts processing', async () => { + // Given - use a handler that takes time and tracks concurrency + let concurrentProcessing = 0 + let maxConcurrency = 0 + const testMessageIds = ['concurrency-1', 'concurrency-2', 'concurrency-3'] + const processedMessages: string[] = [] + + consumer = new PermissionConsumer(testContext.cradle, { + handlers: { + 'permission-added': new KafkaHandlerConfig(PERMISSION_ADDED_SCHEMA, async (message) => { + // Only track messages from this test + if (!testMessageIds.includes(message.value.id)) { + consumer!.addedMessages.push(message) + return + } + + concurrentProcessing++ + maxConcurrency = Math.max(maxConcurrency, concurrentProcessing) + + // Simulate processing time + await new Promise((resolve) => setTimeout(resolve, 30)) + + concurrentProcessing-- + processedMessages.push(message.value.id) + consumer!.addedMessages.push(message) + }), + }, + }) + await consumer.init() + + // When - publish multiple messages + await Promise.all([ + publisher.publish('permission-added', { + id: 'concurrency-1', + type: 'added', + permissions: [], + }), + publisher.publish('permission-added', { + id: 'concurrency-2', + type: 'added', + permissions: [], + }), + publisher.publish('permission-added', { + id: 'concurrency-3', + type: 'added', + permissions: [], + }), + ]) + + // Then - wait for all messages + await consumer.handlerSpy.waitForMessageWithId('concurrency-1', 'consumed') + await consumer.handlerSpy.waitForMessageWithId('concurrency-2', 'consumed') + await consumer.handlerSpy.waitForMessageWithId('concurrency-3', 'consumed') + + // Verify only one message was processed at a time (max concurrency = 1) + expect(maxConcurrency).toBe(1) + expect(processedMessages).toHaveLength(3) + expect(processedMessages).toContain('concurrency-1') + expect(processedMessages).toContain('concurrency-2') + expect(processedMessages).toContain('concurrency-3') + }) + + it('should process messages synchronously across different topics', async () => { + // Given + const processingOrder: string[] = [] + const testMessageIds = ['cross-topic-1', 'cross-topic-2', 'cross-topic-3'] + + consumer = new PermissionConsumer(testContext.cradle, { + handlers: { + 'permission-added': new KafkaHandlerConfig(PERMISSION_ADDED_SCHEMA, async (message) => { + // Only track messages from this test + if (!testMessageIds.includes(message.value.id)) { + consumer!.addedMessages.push(message) + return + } + processingOrder.push(`added-${message.value.id}`) + await new Promise((resolve) => setTimeout(resolve, 20)) + consumer!.addedMessages.push(message) + }), + 'permission-removed': new KafkaHandlerConfig( + PERMISSION_REMOVED_SCHEMA, + async (message) => { + // Only track messages from this test + if (!testMessageIds.includes(message.value.id)) { + consumer!.removedMessages.push(message) + return + } + processingOrder.push(`removed-${message.value.id}`) + await new Promise((resolve) => setTimeout(resolve, 20)) + consumer!.removedMessages.push(message) + }, + ), + }, + }) + await consumer.init() + + // When - publish messages to different topics + await Promise.all([ + publisher.publish('permission-added', { + id: 'cross-topic-1', + type: 'added', + permissions: [], + }), + publisher.publish('permission-removed', { + id: 'cross-topic-2', + type: 'removed', + permissions: [], + }), + publisher.publish('permission-added', { + id: 'cross-topic-3', + type: 'added', + permissions: [], + }), + ]) + + // Then - wait for all messages + await consumer.handlerSpy.waitForMessageWithId('cross-topic-1', 'consumed') + await consumer.handlerSpy.waitForMessageWithId('cross-topic-2', 'consumed') + await consumer.handlerSpy.waitForMessageWithId('cross-topic-3', 'consumed') + + // Verify messages were processed sequentially (one at a time) + // Note: The exact order depends on Kafka's partition assignment, but each should complete before next starts + expect(processingOrder.length).toBe(3) + const testMessages = + consumer.addedMessages.filter((m) => testMessageIds.includes(m.value.id)).length + + consumer.removedMessages.filter((m) => testMessageIds.includes(m.value.id)).length + expect(testMessages).toBe(3) }) }) }) From 207da8574fa3dd90d9352df4ae8f2a66c4778f0b Mon Sep 17 00:00:00 2001 From: Irfan Hodzic Date: Fri, 21 Nov 2025 14:52:13 +0100 Subject: [PATCH 4/4] increasing version to 0.7.7 --- packages/kafka/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/kafka/package.json b/packages/kafka/package.json index e4850886..6ad87275 100644 --- a/packages/kafka/package.json +++ b/packages/kafka/package.json @@ -1,6 +1,6 @@ { "name": "@message-queue-toolkit/kafka", - "version": "0.7.6", + "version": "0.7.7", "engines": { "node": ">= 22.14.0" },