All files / apps/audit-service/src/audit audit-stream.consumer.ts

100% Statements 20/20
100% Branches 10/10
100% Functions 4/4
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 522x           2x 2x 2x     2x     5x     5x 5x       5x           2x     2x       2x   2x 1x     1x 1x           1x      
import {
  Injectable,
  Logger,
  OnApplicationBootstrap,
  OnModuleDestroy,
} from '@nestjs/common';
import { Consumer, Kafka } from 'kafkajs';
import { TOPICS, WorkOrderEvent } from '@app/contracts';
import { AuditIngestService } from './audit-ingest.service';
 
@Injectable()
export class AuditStreamConsumer
  implements OnApplicationBootstrap, OnModuleDestroy
{
  private readonly logger = new Logger(AuditStreamConsumer.name);
  private readonly consumer: Consumer;
 
  constructor(private readonly ingest: AuditIngestService) {
    const kafka = new Kafka({
      clientId: 'audit-service',
      brokers: (process.env.KAFKA_BROKERS ?? 'localhost:9092').split(','),
    });
    this.consumer = kafka.consumer({
      groupId: process.env.AUDIT_CONSUMER_GROUP ?? 'audit-service',
    });
  }
 
  async onApplicationBootstrap(): Promise<void> {
    await this.consumer.connect();
    // fromBeginning is what makes the audit log rebuildable: point a fresh
    // consumer group at the topic and the whole history streams back in.
    await this.consumer.subscribe({
      topic: TOPICS.WORK_ORDER_EVENTS,
      fromBeginning: true,
    });
    await this.consumer.run({
      eachMessage: async ({ message }) => {
        if (!message.value) return;
        const event = JSON.parse(message.value.toString()) as WorkOrderEvent;
        // A thrown error here means the offset is not committed and the
        // message is redelivered; ingest idempotency makes that safe.
        await this.ingest.record(event);
        this.logger.debug(`recorded ${event.type} (${event.eventId})`);
      },
    });
  }
 
  async onModuleDestroy(): Promise<void> {
    await this.consumer.disconnect();
  }
}