Skip to content

Commit 0fde65d

Browse files
authored
feat: update message-broker to version 1.9.0 and add prefetch option to listenOn method (#231)
1 parent bb98511 commit 0fde65d

5 files changed

Lines changed: 45 additions & 7 deletions

File tree

package-lock.json

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

packages/message-broker/package.json

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
{
22
"name": "@user-office-software/duo-message-broker",
3-
"version": "1.8.0",
3+
"version": "1.9.0",
44
"description": "",
55
"author": "SWAP",
66
"license": "ISC",

packages/message-broker/src/consumer.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,15 @@
1-
import { ConsumerCallback } from './index';
1+
import { ConsumerCallback, ListenOnOptions } from './index';
22

33
// This class is used to store the callback function and whether it is registered or not.
44
export class Consumer {
55
callback: ConsumerCallback;
66
registered: boolean;
7+
options: ListenOnOptions;
78

8-
constructor(callback: ConsumerCallback) {
9+
constructor(callback: ConsumerCallback, options: ListenOnOptions = {}) {
910
this.callback = callback;
1011
this.registered = false;
12+
this.options = options;
1113
}
1214

1315
register() {

packages/message-broker/src/index.spec.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ describe('RabbitMQMessageBroker', () => {
2121
bindQueue: jest.fn(),
2222
publish: jest.fn(),
2323
consume: jest.fn().mockResolvedValue({}),
24+
prefetch: jest.fn(),
2425
on: jest.fn(),
2526
};
2627

@@ -109,6 +110,25 @@ describe('RabbitMQMessageBroker', () => {
109110
await broker.listenOn(Queue.SCHEDULING_PROPOSAL, consumer);
110111
expect(mockAmqpChannel.consume).toHaveBeenCalledTimes(2);
111112
});
113+
114+
it('should call channel.prefetch with the given value before consume when prefetch option is set', async () => {
115+
await broker.listenOn(Queue.SCHEDULING_PROPOSAL, consumer, {
116+
prefetch: 10,
117+
});
118+
119+
expect(mockAmqpChannel.prefetch).toHaveBeenCalledWith(10, false);
120+
const prefetchOrder = (mockAmqpChannel.prefetch as jest.Mock).mock
121+
.invocationCallOrder[0];
122+
const consumeOrder = (mockAmqpChannel.consume as jest.Mock).mock
123+
.invocationCallOrder[0];
124+
expect(prefetchOrder).toBeLessThan(consumeOrder);
125+
});
126+
127+
it('should not call channel.prefetch when prefetch option is not set', async () => {
128+
await broker.listenOn(Queue.SCHEDULING_PROPOSAL, consumer);
129+
130+
expect(mockAmqpChannel.prefetch).not.toHaveBeenCalled();
131+
});
112132
});
113133

114134
describe('scheduleReconnect', () => {

packages/message-broker/src/index.ts

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,10 @@ export type ConsumerCallback = (
2424
properties: MessageProperties
2525
) => Promise<void>;
2626

27+
export type ListenOnOptions = {
28+
prefetch?: number;
29+
};
30+
2731
export interface MessageBroker {
2832
sendMessage(queue: Queue, type: string, message: string): Promise<void>;
2933
sendBroadcast(queue: Queue, type: string, message: string): Promise<void>;
@@ -36,7 +40,11 @@ export interface MessageBroker {
3640
queueName: string,
3741
exchangeName: string
3842
): Promise<void>;
39-
listenOn(queue: Queue, cb: ConsumerCallback): Promise<void>;
43+
listenOn(
44+
queue: Queue,
45+
cb: ConsumerCallback,
46+
options?: ListenOnOptions
47+
): Promise<void>;
4048
listenOnBroadcast(cb: ConsumerCallback): void;
4149
}
4250

@@ -110,12 +118,16 @@ export class RabbitMQMessageBroker implements MessageBroker {
110118
}
111119
}
112120

113-
async listenOn(queue: Queue, cb: ConsumerCallback) {
121+
async listenOn(
122+
queue: Queue,
123+
cb: ConsumerCallback,
124+
options: ListenOnOptions = {}
125+
) {
114126
if (!this.queueConsumers.has(queue)) {
115127
this.queueConsumers.set(queue, []);
116128
}
117129

118-
this.queueConsumers.get(queue)?.push(new Consumer(cb));
130+
this.queueConsumers.get(queue)?.push(new Consumer(cb, options));
119131

120132
if (this.channel) {
121133
await this.registerConsumers();
@@ -385,6 +397,10 @@ export class RabbitMQMessageBroker implements MessageBroker {
385397
consumer.register();
386398
}
387399

400+
if (consumer.options.prefetch !== undefined) {
401+
await this.channel.prefetch(consumer.options.prefetch, false);
402+
}
403+
388404
await this.channel
389405
.consume(
390406
queue,

0 commit comments

Comments
 (0)