All files / apps/notifications-service/src/notifications event-retry.handler.ts

100% Statements 19/19
100% Branches 10/10
100% Functions 3/3
100% Lines 17/17

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 633x 3x   3x 3x                         3x 8x   8x               7x 7x   6x   6x 4x     4x         2x     2x                         6x 6x      
import { AmqpConnection } from '@golevelup/nestjs-rabbitmq';
import { Injectable, Logger } from '@nestjs/common';
import type { ConsumeMessage } from 'amqplib';
import { EXCHANGES } from '@app/contracts';
import {
  MAX_ATTEMPTS,
  RETRY_DELAY_MS,
  retryQueue,
} from './messaging-resilience';
 
/**
 * Wraps a message handler with retry + dead-letter semantics. The original
 * message is always acked; on failure a copy is republished either to the
 * retry queue (delayed redelivery) or, once attempts are exhausted, to the
 * dead-letter exchange. Attempt count travels in the x-attempt header.
 */
@Injectable()
export class EventRetryHandler {
  private readonly logger = new Logger(EventRetryHandler.name);
 
  constructor(private readonly amqp: AmqpConnection) {}
 
  async handle(
    queue: string,
    event: object,
    message: ConsumeMessage,
    handler: () => Promise<void>,
  ): Promise<void> {
    try {
      await handler();
    } catch (error) {
      const attempt = this.attemptOf(message);
 
      if (attempt < MAX_ATTEMPTS) {
        this.logger.warn(
          `attempt ${attempt}/${MAX_ATTEMPTS} failed on ${queue}, retrying in ${RETRY_DELAY_MS}ms`,
        );
        await this.amqp.publish('', retryQueue(queue), event, {
          persistent: true,
          headers: { 'x-attempt': attempt + 1 },
        });
      } else {
        this.logger.error(
          `attempt ${attempt}/${MAX_ATTEMPTS} failed on ${queue}, dead-lettering`,
        );
        await this.amqp.publish(EXCHANGES.DEAD_LETTER, queue, event, {
          persistent: true,
          headers: {
            'x-attempt': attempt,
            'x-last-error':
              error instanceof Error ? error.message : String(error),
          },
        });
      }
    }
  }
 
  private attemptOf(message: ConsumeMessage): number {
    const raw: unknown = message.properties.headers?.['x-attempt'];
    return typeof raw === 'number' ? raw : 1;
  }
}