From 8e77dc27a66b4ad112d73c275ab645d661e0ef6b Mon Sep 17 00:00:00 2001 From: sebastian Date: Tue, 16 Jun 2026 14:34:19 +0200 Subject: [PATCH 1/2] [messaging] fix: channel wrapper avoid override --- .../messaging-rabbitmq-extension/package.json | 2 +- .../consumer/rabbitmq-messaging.consumer.ts | 28 +++++++++++-------- 2 files changed, 17 insertions(+), 13 deletions(-) diff --git a/packages/messaging-rabbitmq-extension/package.json b/packages/messaging-rabbitmq-extension/package.json index 89d2572..a91e470 100644 --- a/packages/messaging-rabbitmq-extension/package.json +++ b/packages/messaging-rabbitmq-extension/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging-rabbitmq-extension", - "version": "4.2.0", + "version": "4.2.1", "description": "Extension to handle messages and dispatch them over AMQP protocol", "author": "Sebastian Iwanczyszyn", "license": "MIT", diff --git a/packages/messaging-rabbitmq-extension/src/consumer/rabbitmq-messaging.consumer.ts b/packages/messaging-rabbitmq-extension/src/consumer/rabbitmq-messaging.consumer.ts index e5abddb..2c65f1d 100644 --- a/packages/messaging-rabbitmq-extension/src/consumer/rabbitmq-messaging.consumer.ts +++ b/packages/messaging-rabbitmq-extension/src/consumer/rabbitmq-messaging.consumer.ts @@ -19,8 +19,8 @@ import { MessageDeadLetterVisitor } from './message-dead-letter.visitor'; export class RabbitmqMessagingConsumer implements IMessagingConsumer, OnModuleDestroy { - private channel?: AmqpChannel = undefined; - private amqpChannel: ChannelWrapper; + private readonly channelWrappers = new WeakMap(); + private readonly channels = new Set(); constructor( private readonly rabbitMqMigrator: RabbitmqMigrator, @@ -32,7 +32,7 @@ export class RabbitmqMessagingConsumer dispatcher: ConsumerMessageBus, channel: AmqpChannel, ): Promise { - this.channel = channel; + this.channels.add(channel); await this.rabbitMqMigrator.run(channel); if (!channel.connection) { @@ -43,10 +43,10 @@ export class RabbitmqMessagingConsumer const channelWrapper = channel.createChannelWrapper(); await channelWrapper.waitForConnect(); - this.amqpChannel = channelWrapper; + this.channelWrappers.set(channel, channelWrapper); await channelWrapper.addSetup(async (rawChannel: Channel) => { - await rawChannel.prefetch(this.channel.config.qos, false); + await rawChannel.prefetch(channel.config.qos, false); return rawChannel.consume( channel.config.queue, async (msg: ConsumeMessage | null) => { @@ -89,7 +89,9 @@ export class RabbitmqMessagingConsumer errored: ConsumerDispatchedMessageError, channel: AmqpChannel, ): Promise { - if (!this.amqpChannel) { + const channelWrapper = this.channelWrappers.get(channel); + + if (!channelWrapper) { return; } @@ -104,25 +106,27 @@ export class RabbitmqMessagingConsumer return this.messageRetrier.retryMessage( errored, channel, - this.amqpChannel, + channelWrapper, currentRetryCount, ); } } - if (channel.config.deadLetterQueueFeature && this.amqpChannel) { + if (channel.config.deadLetterQueueFeature) { return this.messageDeadLetter.sendToDeadLetter( errored, channel, - this.amqpChannel, + channelWrapper, ); } } async onModuleDestroy(): Promise { - if (this.channel?.connection) { - await this.channel.connection.close(); + for (const channel of this.channels) { + if (channel.connection) { + await channel.connection.close(); + } } - this.channel = undefined; + this.channels.clear(); } } From 7b1203e9340f9b6a202d9563a94d3af060f0d41e Mon Sep 17 00:00:00 2001 From: sebastian Date: Tue, 16 Jun 2026 14:39:27 +0200 Subject: [PATCH 2/2] [messaging] fix: channel wrapper avoid override --- .../rabbitmq-messaging.consumer.spec.ts | 33 ++++++++++++------- 1 file changed, 21 insertions(+), 12 deletions(-) diff --git a/packages/messaging-rabbitmq-extension/test/unit/consumer/rabbitmq-messaging.consumer.spec.ts b/packages/messaging-rabbitmq-extension/test/unit/consumer/rabbitmq-messaging.consumer.spec.ts index b4f6c31..d20716b 100644 --- a/packages/messaging-rabbitmq-extension/test/unit/consumer/rabbitmq-messaging.consumer.spec.ts +++ b/packages/messaging-rabbitmq-extension/test/unit/consumer/rabbitmq-messaging.consumer.spec.ts @@ -120,7 +120,7 @@ describe('RabbitmqMessagingConsumer', () => { } as unknown as jest.Mocked; const setupFn = (wrapper.addSetup as jest.Mock).mock.calls[0][0]; await setupFn(rawChannel); - expect(rawChannel.prefetch).toHaveBeenCalledWith(10); + expect(rawChannel.prefetch).toHaveBeenCalledWith(10, false); const consumeHandler = (rawChannel.consume as jest.Mock).mock .calls[0][1] as (msg: ConsumeMessage | null) => Promise; @@ -230,13 +230,16 @@ describe('RabbitmqMessagingConsumer', () => { } as any; amqpChannel = {} as any; - - (consumer as any).channel = channel; - (consumer as any).amqpChannel = amqpChannel; + (consumer as any).channelWrappers.set(channel, amqpChannel); }); it('should return when no amqpChannel is available', async () => { - (consumer as any).amqpChannel = undefined; + const unregisteredChannel = { + config: { + retryMessage: 3, + deadLetterQueueFeature: true, + }, + } as any; const errored: ConsumerDispatchedMessageError = { dispatchedConsumerMessage: { @@ -244,7 +247,7 @@ describe('RabbitmqMessagingConsumer', () => { }, } as any; - const result = await consumer.onError(errored, channel); + const result = await consumer.onError(errored, unregisteredChannel); expect(result).toBeUndefined(); expect(mockRetrier.retryMessage).not.toHaveBeenCalled(); @@ -311,6 +314,7 @@ describe('RabbitmqMessagingConsumer', () => { deadLetterQueueFeature: true, }, } as any; + (consumer as any).channelWrappers.set(channelWithoutRetry, amqpChannel); const errored: ConsumerDispatchedMessageError = { dispatchedConsumerMessage: { @@ -335,6 +339,10 @@ describe('RabbitmqMessagingConsumer', () => { deadLetterQueueFeature: false, }, } as any; + (consumer as any).channelWrappers.set( + channelWithDisabledDeadLetter, + amqpChannel, + ); const errored: ConsumerDispatchedMessageError = { dispatchedConsumerMessage: { @@ -367,26 +375,27 @@ describe('RabbitmqMessagingConsumer', () => { }); describe('onModuleDestroy', () => { - it('should close connection and clear channel reference', async () => { + it('should close all registered channel connections', async () => { const close = jest.fn().mockResolvedValue(undefined); - (consumer as any).channel = { + const channel = { connection: { close, }, }; + (consumer as any).channels.add(channel); await consumer.onModuleDestroy(); expect(close).toHaveBeenCalledTimes(1); - expect((consumer as any).channel).toBeUndefined(); + expect((consumer as any).channels.size).toBe(0); }); - it('should only clear channel reference when connection does not exist', async () => { - (consumer as any).channel = { connection: undefined }; + it('should ignore channels without connection and still clear registry', async () => { + (consumer as any).channels.add({ connection: undefined }); await consumer.onModuleDestroy(); - expect((consumer as any).channel).toBeUndefined(); + expect((consumer as any).channels.size).toBe(0); }); }); });