Skip to content

Commit adcb88b

Browse files
committed
[CHORE]: Enhance DLQ handling with error logging and channel state checks
1 parent 6537070 commit adcb88b

1 file changed

Lines changed: 38 additions & 3 deletions

File tree

src/index.js

Lines changed: 38 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -247,10 +247,41 @@ const drainDLQ = async (channel) => {
247247
const connection = await amqp.connect(process.env.AMQ_URL, {
248248
heartbeat: 30
249249
})
250+
251+
let channelOpen = true;
252+
253+
connection.on('error', (err) => {
254+
log('error', `AMQP connection error: ${err.message}`);
255+
});
256+
connection.on('close', () => {
257+
log('error', 'AMQP connection closed unexpectedly. Exiting for restart...');
258+
channelOpen = false;
259+
clearInterval(drainInterval);
260+
process.exit(1);
261+
});
262+
250263
const channel = await connection.createChannel();
264+
265+
channel.on('error', (err) => {
266+
log('error', `AMQP channel error: ${err.message}`);
267+
channelOpen = false;
268+
});
269+
channel.on('close', () => {
270+
log('error', 'AMQP channel closed. Exiting for restart...');
271+
channelOpen = false;
272+
clearInterval(drainInterval);
273+
process.exit(1);
274+
});
275+
251276
await channel.assertQueue(queueName, { durable: false });
252277
await channel.assertQueue(dlqName, { durable: true });
253-
const drainInterval = setInterval(() => drainDLQ(channel), DLQ_DRAIN_INTERVAL_MS);
278+
const drainInterval = setInterval(() => {
279+
if (!channelOpen) {
280+
log('warn', 'Skipping DLQ drain: channel is not open');
281+
return;
282+
}
283+
drainDLQ(channel);
284+
}, DLQ_DRAIN_INTERVAL_MS);
254285
log('info', `DLQ drain scheduled every ${DLQ_DRAIN_INTERVAL_MS / 1000}s`);
255286
await channel.consume(queueName, async (message) => {
256287
let queueMsg;
@@ -291,8 +322,12 @@ const drainDLQ = async (channel) => {
291322

292323
} catch (error) {
293324
log('error', error.message);
294-
sendToDLQ(channel, message, error.message);
295-
channel.ack(message);
325+
if (channelOpen) {
326+
sendToDLQ(channel, message, error.message);
327+
channel.ack(message);
328+
} else {
329+
log('warn', 'Channel closed during message processing; message will be requeued on restart');
330+
}
296331
}
297332

298333
})

0 commit comments

Comments
 (0)