From cec3f1e6b1da8c6b93420c993cc69d2e5a47a55f Mon Sep 17 00:00:00 2001 From: jekabs-karklins Date: Wed, 22 Jul 2026 13:59:43 +0200 Subject: [PATCH] feat(queue): add event type to log messages in QueueConsumer --- src/queue/consumers/QueueConsumer.spec.ts | 21 ++++++++++++++++----- src/queue/consumers/QueueConsumer.ts | 3 ++- 2 files changed, 18 insertions(+), 6 deletions(-) diff --git a/src/queue/consumers/QueueConsumer.spec.ts b/src/queue/consumers/QueueConsumer.spec.ts index 4aa2945c..fe9eff79 100644 --- a/src/queue/consumers/QueueConsumer.spec.ts +++ b/src/queue/consumers/QueueConsumer.spec.ts @@ -82,33 +82,44 @@ describe('QueueConsumer', () => { }); it('should call onMessage when a message is received', async () => { + const eventType = 'PROPOSAL_CREATED'; + const message = { proposalPk: 1234 }; + await queueConsumer.start(); await messageBrokerMock.listenOn.mock.calls[0][1]( - 'testMessage', - {}, + eventType, + message, {} as any ); expect(logger.logInfo).toHaveBeenCalledWith('Received message on queue', { queueName: 'testQueue', + eventType, }); expect(logger.logException).not.toHaveBeenCalled(); - expect(queueConsumer.onMessage).toHaveBeenCalledWith('testMessage', {}, {}); + expect(queueConsumer.onMessage).toHaveBeenCalledWith( + eventType, + message, + {} + ); }); it('should log error if onMessage throws and rethrow the error', async () => { + const eventType = 'PROPOSAL_CREATED'; + const message = { proposalPk: 1234 }; const error = new Error('Test error'); queueConsumer.onMessage = jest.fn().mockRejectedValue(error); await queueConsumer.start(); await expect( - messageBrokerMock.listenOn.mock.calls[0][1]('testMessage', {}, {} as any) + messageBrokerMock.listenOn.mock.calls[0][1](eventType, message, {} as any) ).rejects.toThrow(error); expect(logger.logInfo).toHaveBeenCalledWith('Received message on queue', { queueName: 'testQueue', + eventType, }); expect(logger.logException).toHaveBeenCalledWith( 'Error while handling QueueConsumer callback: ', @@ -116,7 +127,7 @@ describe('QueueConsumer', () => { error: error.message, queue: 'testQueue', consumer: 'TestQueueConsumer', - args: ['testMessage', {}, {}], + args: [eventType, message, {}], } ); }); diff --git a/src/queue/consumers/QueueConsumer.ts b/src/queue/consumers/QueueConsumer.ts index abff40b2..85c82888 100644 --- a/src/queue/consumers/QueueConsumer.ts +++ b/src/queue/consumers/QueueConsumer.ts @@ -65,7 +65,8 @@ export abstract class QueueConsumer { this.messageBroker.listenOn( queueName as Queue, async (...args) => { - logger.logInfo('Received message on queue', { queueName }); + const [eventType] = args; + logger.logInfo('Received message on queue', { queueName, eventType }); // Start tracking processing time const endTimer = processingDurationHistogram.startTimer({