All files / apps/work-orders-service/src/messaging outbox-relay.ts

100% Statements 33/33
100% Branches 20/20
100% Functions 6/6
100% Lines 28/28

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 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 852x 2x           2x 2x 2x 2x   2x                             2x 10x 10x 10x     10x 10x 10x       3x 4x       6x       9x 8x 8x 8x 8x                   7x 2x       1x 1x           2x           8x        
import { AmqpConnection } from '@golevelup/nestjs-rabbitmq';
import {
  Injectable,
  Logger,
  OnApplicationBootstrap,
  OnModuleDestroy,
} from '@nestjs/common';
import { DataSource } from 'typeorm';
import { EXCHANGES } from '@app/contracts';
import { OutboxEvent } from './outbox-event.entity';
import { WorkOrderAuditProducer } from './work-order-audit.producer';
 
const BATCH_SIZE = 20;
 
/**
 * The publishing half of the outbox: polls the unpublished tail and forwards
 * each event to RabbitMQ and the Kafka audit stream. Rows are claimed with
 * FOR UPDATE SKIP LOCKED, so multiple service instances can run relays
 * without publishing the same row twice (a crashed instance's rows are simply
 * picked up by the next tick elsewhere).
 *
 * Failure semantics: a broker error aborts the tick and rolls the batch's
 * published_at stamps back — the rows are retried next tick. That makes
 * delivery at-least-once, which is exactly what every consumer in the system
 * is already built for (idempotency via event_id).
 */
@Injectable()
export class OutboxRelay implements OnApplicationBootstrap, OnModuleDestroy {
  private readonly logger = new Logger(OutboxRelay.name);
  private timer: NodeJS.Timeout | null = null;
  private draining = false;
 
  constructor(
    private readonly dataSource: DataSource,
    private readonly amqp: AmqpConnection,
    private readonly audit: WorkOrderAuditProducer,
  ) {}
 
  onApplicationBootstrap(): void {
    const intervalMs = parseInt(process.env.OUTBOX_POLL_MS ?? '500', 10);
    this.timer = setInterval(() => void this.drain(), intervalMs);
  }
 
  onModuleDestroy(): void {
    if (this.timer) clearInterval(this.timer);
  }
 
  async drain(): Promise<void> {
    if (this.draining) return; // a slow tick must not overlap the next one
    this.draining = true;
    try {
      await this.dataSource.transaction(async (manager) => {
        const rows = await manager
          .getRepository(OutboxEvent)
          .createQueryBuilder('outbox')
          .setLock('pessimistic_write')
          .setOnLocked('skip_locked')
          .where('outbox.publishedAt IS NULL')
          .orderBy('outbox.id', 'ASC')
          .take(BATCH_SIZE)
          .getMany();
 
        for (const row of rows) {
          await this.amqp.publish(EXCHANGES.EVENTS, row.type, row.payload, {
            persistent: true,
            messageId: row.eventId,
          });
          await this.audit.record(row.payload);
          await manager.update(OutboxEvent, row.id, {
            publishedAt: new Date(),
          });
        }
      });
    } catch (error) {
      this.logger.warn(
        `outbox drain failed, will retry: ${
          error instanceof Error ? error.message : String(error)
        }`,
      );
    } finally {
      this.draining = false;
    }
  }
}