Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 16 additions & 5 deletions src/queue/consumers/QueueConsumer.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -82,41 +82,52 @@ 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: ',
{
error: error.message,
queue: 'testQueue',
consumer: 'TestQueueConsumer',
args: ['testMessage', {}, {}],
args: [eventType, message, {}],
}
);
});
Expand Down
3 changes: 2 additions & 1 deletion src/queue/consumers/QueueConsumer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down
Loading