diff --git a/packages/rabbitmq/src/options.ts b/packages/rabbitmq/src/options.ts index 49916f1..2c453ca 100644 --- a/packages/rabbitmq/src/options.ts +++ b/packages/rabbitmq/src/options.ts @@ -49,6 +49,8 @@ export interface RabbitMQTransportOptions { deadLetterUnhandled?: boolean; queueArguments?: Record; retryQueueArguments?: Record; + errorQueueArguments?: Record; + auditQueueArguments?: Record; }; } @@ -70,6 +72,8 @@ export interface ResolvedConsumerOptions { readonly deadLetterUnhandled: boolean; readonly queueArguments: Readonly>; readonly retryQueueArguments: Readonly>; + readonly errorQueueArguments: Readonly>; + readonly auditQueueArguments: Readonly>; } export function resolveProducerOptions(opts: RabbitMQTransportOptions): ResolvedProducerOptions { @@ -94,5 +98,7 @@ export function resolveConsumerOptions(opts: RabbitMQTransportOptions): Resolved deadLetterUnhandled: c.deadLetterUnhandled ?? false, queueArguments: c.queueArguments ?? {}, retryQueueArguments: c.retryQueueArguments ?? {}, + errorQueueArguments: c.errorQueueArguments ?? {}, + auditQueueArguments: c.auditQueueArguments ?? {}, }; } diff --git a/packages/rabbitmq/src/topology.ts b/packages/rabbitmq/src/topology.ts index 251ee5f..229b3ee 100644 --- a/packages/rabbitmq/src/topology.ts +++ b/packages/rabbitmq/src/topology.ts @@ -84,12 +84,20 @@ export function buildConsumerTopology( if (opts.errorQueue !== null) { exchanges.push({ exchange: opts.errorQueue, type: 'direct', durable: false }); - queues.push({ queue: opts.errorQueue, durable: true }); + queues.push({ + queue: opts.errorQueue, + durable: true, + arguments: { ...opts.errorQueueArguments }, + }); queueBindings.push({ exchange: opts.errorQueue, queue: opts.errorQueue, routingKey: '' }); } if (opts.auditEnabled) { exchanges.push({ exchange: opts.auditQueue, type: 'direct', durable: false }); - queues.push({ queue: opts.auditQueue, durable: true }); + queues.push({ + queue: opts.auditQueue, + durable: true, + arguments: { ...opts.auditQueueArguments }, + }); queueBindings.push({ exchange: opts.auditQueue, queue: opts.auditQueue, routingKey: '' }); } diff --git a/packages/rabbitmq/test/unit/options.test.ts b/packages/rabbitmq/test/unit/options.test.ts index d5e999d..0a49b34 100644 --- a/packages/rabbitmq/test/unit/options.test.ts +++ b/packages/rabbitmq/test/unit/options.test.ts @@ -39,6 +39,8 @@ describe('resolveConsumerOptions', () => { expect(resolved.deadLetterUnhandled).toBe(false); expect(resolved.queueArguments).toEqual({}); expect(resolved.retryQueueArguments).toEqual({}); + expect(resolved.errorQueueArguments).toEqual({}); + expect(resolved.auditQueueArguments).toEqual({}); }); it('caller overrides win', () => { @@ -51,6 +53,8 @@ describe('resolveConsumerOptions', () => { errorQueue: 'my-errors', auditEnabled: true, queueArguments: { 'x-max-priority': 10 }, + errorQueueArguments: { 'x-queue-type': 'quorum' }, + auditQueueArguments: { 'x-queue-type': 'quorum' }, }, }; const resolved = resolveConsumerOptions(opts); @@ -60,6 +64,8 @@ describe('resolveConsumerOptions', () => { expect(resolved.errorQueue).toBe('my-errors'); expect(resolved.auditEnabled).toBe(true); expect(resolved.queueArguments).toEqual({ 'x-max-priority': 10 }); + expect(resolved.errorQueueArguments).toEqual({ 'x-queue-type': 'quorum' }); + expect(resolved.auditQueueArguments).toEqual({ 'x-queue-type': 'quorum' }); }); it('errorQueue: null opts out of error-queue routing', () => { diff --git a/packages/rabbitmq/test/unit/topology.test.ts b/packages/rabbitmq/test/unit/topology.test.ts index df9c7fb..edf43cb 100644 --- a/packages/rabbitmq/test/unit/topology.test.ts +++ b/packages/rabbitmq/test/unit/topology.test.ts @@ -60,7 +60,7 @@ describe('retry/error/audit topology (master parity)', () => { type: 'direct', durable: false, }); - expect(topo.queues).toContainEqual({ queue: 'errors', durable: true }); + expect(topo.queues).toContainEqual({ queue: 'errors', durable: true, arguments: {} }); expect(topo.queueBindings).toContainEqual({ exchange: 'errors', queue: 'errors', @@ -143,6 +143,32 @@ describe('buildConsumerTopology (general)', () => { const main = t.queues.find((q) => q.queue === 'q-self'); expect(main?.arguments?.['x-max-priority']).toBe(10); }); + + it('merges caller errorQueueArguments into the error queue', () => { + const t = buildConsumerTopology( + 'q-self', + [], + resolveConsumerOptions({ + url: '', + consumer: { errorQueueArguments: { 'x-queue-type': 'quorum' } }, + }), + ); + const error = t.queues.find((q) => q.queue === 'errors'); + expect(error?.arguments?.['x-queue-type']).toBe('quorum'); + }); + + it('merges caller auditQueueArguments into the audit queue', () => { + const t = buildConsumerTopology( + 'q-self', + [], + resolveConsumerOptions({ + url: '', + consumer: { auditEnabled: true, auditQueueArguments: { 'x-queue-type': 'quorum' } }, + }), + ); + const audit = t.queues.find((q) => q.queue === 'audit'); + expect(audit?.arguments?.['x-queue-type']).toBe('quorum'); + }); }); describe('exchangeNameForType', () => {