@miguelcazares/invoixup-events
v1.0.2
Published
Contratos de eventos de Kafka de invoixup (topics, payloads, DTOs) y toolkit NestJS para consumidores (reintentos + dead letter).
Maintainers
Readme
📦 @miguelcazares/invoixup-events
Fuente de verdad única de la mensajería Kafka de invoixup: nombres de topics, contratos de payload, DTOs de validación y el toolkit de NestJS que usan los consumidores (reintentos en proceso + dead letter).
Cualquier microservicio que quiera producir o consumir un evento importa las constantes y los tipos desde aquí, en vez de duplicarlos.
🚀 Instalación
pnpm add @miguelcazares/invoixup-eventsEl paquete tiene dos entradas:
| Import | Contenido | Peers que necesita |
|---|---|---|
| @miguelcazares/invoixup-events | topics, payloads, DTOs, headers del DLQ | class-validator, class-transformer, reflect-metadata |
| @miguelcazares/invoixup-events/nestjs | KafkaConsumerModule, interceptor, DLQ, filtro rpc | además: @nestjs/common, @nestjs/microservices, kafkajs, rxjs, @sentry/nestjs |
Los peers del toolkit son opcionales: un servicio que solo consume los contratos no necesita instalarlos.
📇 Contratos
import {
KAFKA_TOPICS,
ALL_TOPICS,
dlqTopic,
KAFKA_CONSUMER_GROUPS,
serverGroupId,
DLQ_HEADERS,
UserRegisteredEventV1,
UserRegisteredEventDto,
EventPayloadMap,
} from '@miguelcazares/invoixup-events';
KAFKA_TOPICS.AUTH_USER_REGISTERED; // 'auth.user.registered'
dlqTopic(KAFKA_TOPICS.AUTH_USER_REGISTERED); // 'auth.user.registered.dlq'
serverGroupId(KAFKA_CONSUMER_GROUPS.MS_PAYMENTS); // 'ms-payments-server'Topics
| Topic | Productor | Consumidor |
|---|---|---|
| auth.user.registered | ms-auths | ms-payments |
| auth.user.registered.dlq | ms-payments | inspección manual |
ALL_TOPICS es la lista que invoixup-kafka tiene que crear en el broker.
Mantener la variable TOPICS de su docker-compose.yml sincronizada con ella.
Producir (ms-auths)
El código de negocio encola en el outbox transaccional; solo necesita el nombre del topic y el tipo del payload:
import { KAFKA_TOPICS, UserRegisteredEventV1 } from '@miguelcazares/invoixup-events';
await this.outboxService.enqueue(
manager,
KAFKA_TOPICS.AUTH_USER_REGISTERED,
{
userId, businessId, planId, email, name,
occurredAt: new Date().toISOString(),
} satisfies UserRegisteredEventV1,
String(businessId), // clave de partición
);Consumir (ms-payments)
import { ValidationPipe } from '@nestjs/common';
import { EventPattern, Payload, Ctx, KafkaContext } from '@nestjs/microservices';
import {
KAFKA_TOPICS,
UserRegisteredEventDto,
} from '@miguelcazares/invoixup-events';
import {
KafkaRetryDlqInterceptor,
KafkaExceptionFilter,
} from '@miguelcazares/invoixup-events/nestjs';
@Controller()
@UseInterceptors(KafkaRetryDlqInterceptor)
@UseFilters(KafkaExceptionFilter)
export class SubscriptionsKafkaController {
@EventPattern(KAFKA_TOPICS.AUTH_USER_REGISTERED)
async handleUserRegistered(
// whitelist SIN forbidNonWhitelisted: un campo nuevo del productor se
// ignora en vez de tumbar el consumidor -> el orden de despliegue deja de
// importar para cambios aditivos del contrato.
@Payload(new ValidationPipe({ whitelist: true, transform: true }))
event: UserRegisteredEventDto,
@Ctx() context: KafkaContext,
): Promise<void> {
await this.subscriptionsService.ensureSubscription(event);
}
}🧰 Toolkit NestJS (/nestjs)
KafkaConsumerModule
Registra DeadLetterService y KafkaRetryDlqInterceptor globalmente.
import { KafkaConsumerModule } from '@miguelcazares/invoixup-events/nestjs';
@Module({
imports: [
KafkaConsumerModule.forRootAsync({
inject: [ConfigService],
useFactory: (config: ConfigService) => ({
clientId: config.get('kafka.clientId'),
brokers: config.get('kafka.brokers'),
dlqSuffix: config.get('kafka.dlqSuffix'), // opcional, default '.dlq'
retryAttempts: config.get('kafka.retryAttempts'), // opcional, default 3
retryBackoffMs: config.get('kafka.retryBackoffMs'), // opcional, default 1000
}),
}),
],
})
export class AppModule {}KafkaRetryDlqInterceptor
- Ejecuta el handler.
- Error transitorio → reintenta
retryAttemptsveces con backoff lineal, en proceso (la partición espera; el total debe caber en elsessionTimeoutdel consumer, 30s por defecto). - Error permanente (
HttpException4xx) → sin reintentos, directo al DLQ. - Agotados los intentos → publica en
<topic>.dlq, emite un valor para que el offset avance, y la partición sigue. - Si la publicación al DLQ falla → deja escalar el error (mejor re-bloquearse que perder el mensaje).
DeadLetterService
Producer de kafkajs directo. Reenvía key/value/headers originales y añade el
contexto del fallo en headers x-dlq-* (ver DLQ_HEADERS). Conecta bajo
demanda en el primer publish().
KafkaExceptionFilter
Filtro a nivel de controller para contexto rpc: loguea topic/partition/offset y
reporta a Sentry sin que el GlobalExceptionFilter HTTP-only enmascare el error.
isHttpContext(context)
Helper para guards e interceptors globales que leen request: devuelve false
en un handler de Kafka, donde no hay request.
🔀 Versionado del contrato
- Las interfaces llevan sufijo de versión (
UserRegisteredEventV1); el nombre del topic no cambia.UserRegisteredEventes alias a la versión vigente. - Cambio aditivo (campo opcional nuevo): sube la minor del paquete. Con el
ValidationPiperelajado del consumidor, cualquier orden de despliegue. - Cambio incompatible (renombrar/quitar campo, cambiar tipo): topic nuevo
<topic>.v2+UserRegisteredEventV2. Los consumidores migran y luego se retira el v1. Nunca se rompe un topic en vivo. - Registrar cada versión en
CHANGELOG.md.
🛠️ Desarrollo
pnpm install
pnpm run build # tsc -> dist/ (entradas . y ./nestjs)Publicación: npm publish (el paquete es público bajo el scope
@miguelcazares).
