Arquitecto de Código · CultivaMetrics API

Architecture: Sistema de Notificaciones en Tiempo Real

CultivaMetrics API · Node.js 20 + TypeScript + Prisma + SSE · Generado por CULTIVA IA
Node.js 20 TypeScript 5 Express 4 Prisma ORM SSE EventBus PostgreSQL 16 Vitest Zod
7
Archivos nuevos
4
Modificaciones
5
Interfaces clave
6
Fases de build
🏗️
Decisiones de Diseño
Razonamiento detrás de cada elección arquitectónica
D1 · Patrón módulo existente
Crear src/modules/notifications/ siguiendo la estructura controller/service/repository/types del resto del proyecto.
El codebase ya usa este patrón de forma consistente en campaigns y users. Replicarlo evita fricción cognitiva y mantiene la navegación predecible para el equipo.
D2 · EventBus para desacoplamiento
El CampaignService emite eventos al EventBus existente; el NotificationService es subscriber, nunca se llama directamente.
Cero acoplamiento entre módulos. Futuras fuentes (comentarios, presupuesto) solo añaden un EventBus.emit sin tocar el módulo de notificaciones.
D3 · SSE en lugar de WebSocket
Endpoint GET /notifications/stream con text/event-stream. Conexión unidireccional servidor → cliente.
Las notificaciones son push unidireccional. SSE es HTTP nativo, funciona sin librerías extra, no necesita upgrade de protocolo y el frontend puede usar la API EventSource estándar.
D4 · Cursor-based pagination
El historial de notificaciones pagina con cursor: notificationId + limit (máx 50).
Consistente con el estilo de la API. El offset-based es inestable bajo inserciones concurrentes; el cursor evita duplicados o saltos en el feed en tiempo real.
D5 · SseConnectionManager en lib/
Singleton src/lib/sse.ts gestiona el mapa de conexiones SSE activas por userId.
Centraliza la gestión de conexiones fuera del módulo. Si en el futuro se migra a Redis Pub/Sub solo se modifica este fichero, sin tocar el módulo de notificaciones.
D6 · Zod schema en schemas/
Validación de entrada (markAsRead, query params del historial) con Zod schemas en src/schemas/notification.schema.ts.
Todos los controllers del proyecto usan Zod. Añadir el schema en la carpeta schemas/ mantiene la convención y permite reutilizarlo en tests.
📁
Archivos a Crear
7 ficheros nuevos
Archivo Propósito Prioridad
src/modules/notifications/types.ts NEW
Tipos base del módulo
Interfaces Notification, NotificationType, NotificationPayload. DTO de creación y respuesta de lista. ● Alta
src/modules/notifications/repository.ts NEW
Capa de datos (Prisma)
CRUD sobre tabla Notification. Métodos: create, findByUser (cursor), markRead, countUnread. ● Alta
src/modules/notifications/service.ts NEW
Lógica de negocio
Suscribe al EventBus (campaign.status_changed). Filtra usuarios elegibles, persiste la notificación y envía por SSE. ● Alta
src/modules/notifications/controller.ts NEW
Express router + SSE endpoint
Rutas: GET /stream (SSE), GET / (historial paginado), PATCH /:id/read, GET /unread-count. ● Media
src/schemas/notification.schema.ts NEW
Validación Zod
Schemas Zod para ListNotificationsQuery (cursor?, limit?) y MarkReadParams. ● Media
src/lib/sse.ts NEW
SSE connection manager
Singleton con Map<userId, Response[]>. Métodos: register, unregister, send(userId, event). ● Alta
src/__tests__/notifications.test.ts NEW
Suite de tests (Vitest)
Tests unitarios del service (mock EventBus + mock repo) y tests de integración del controller (supertest). ● Normal
✏️
Archivos a Modificar
4 ficheros existentes — cambios mínimos e invasivos
⚠️
Principio de mínima invasión: todos los cambios en archivos existentes son adiciones (no modificaciones destructivas). Se preserva el comportamiento actual de cada módulo.
Archivo Cambios requeridos Prioridad
src/modules/campaigns/service.ts MOD
Emisión de eventos de estado
Añadir EventBus.emit('campaign.status_changed', { campaignId, oldStatus, newStatus, assignedUserIds }) en el método que actualiza el estado de la campaña. Sin otros cambios. ● Alta
prisma/schema.prisma MOD
Nuevo modelo Notification
Añadir modelo Notification con campos: id, userId, campaignId?, type, payload (Json), read, createdAt. Índices en (userId, read) y (userId, createdAt DESC). ● Alta
src/lib/events.ts MOD
Tipado del EventBus
Añadir type CampaignStatusChangedEvent al mapa de eventos del EventBus para tipado estricto. Sin cambios en la implementación runtime. ● Media
src/app.ts MOD
Registro de rutas y bootstrap
Montar el router de notificaciones en /api/notifications. Importar e inicializar el NotificationService (para que suscriba al EventBus al arrancar). ● Media
🔌
Interfaces Clave
Contratos TypeScript del módulo
Notification
types.ts
interface Notification {
  id:         string
  userId:     string
  campaignId?: string
  type:       NotificationType
  payload:    Record<string, unknown>
  read:       boolean
  createdAt:  Date
}

type NotificationType =
  | 'campaign.status_changed'
  // futuras: 'comment.added', 'budget.exceeded'
CampaignStatusChangedEvent
lib/events.ts
interface CampaignStatusChangedEvent {
  campaignId:      string
  oldStatus:       CampaignStatus
  newStatus:       CampaignStatus
  assignedUserIds: string[]
  changedAt:       Date
}

// EventBus map actualizado:
type EventMap = {
  'campaign.status_changed':
    CampaignStatusChangedEvent
}
INotificationRepository
repository.ts
interface INotificationRepository {
  create(data: CreateNotificationDto):
    Promise<Notification>

  findByUser(
    userId: string,
    cursor?: string,
    limit: number
  ): Promise<PaginatedResult<Notification>>

  markRead(id: string, userId: string):
    Promise<Notification>

  countUnread(userId: string):
    Promise<number>
}
SseConnectionManager
lib/sse.ts
class SseConnectionManager {
  // Map<userId → Response[]>
  private connections: Map<string, Response[]>

  register(
    userId: string,
    res: Response
  ): void

  unregister(userId: string, res: Response): void

  send(
    userId: string,
    event: SseEvent
  ): void
}
Prisma Model — Notification
schema.prisma
model Notification {
  id         String   @id @default(cuid())
  userId     String
  campaignId String?
  type       String   // NotificationType enum a nivel app
  payload    Json
  read       Boolean  @default(false)
  createdAt  DateTime @default(now())

  user     User     @relation(fields: [userId], references: [id])
  campaign Campaign? @relation(fields: [campaignId], references: [id])

  @@index([userId, read])
  @@index([userId, createdAt(sort: Desc)])
}
🔄
Flujo de Datos
Desde el cambio de estado de campaña hasta el cliente React
Camino de escritura (push)
PATCH /campaigns/:id
status update
CampaignService
.updateStatus()
EventBus.emit
'campaign.status_changed'
NotificationService
.onStatusChanged()
NotificationRepo
.create() × N users
SseConnectionManager
.send(userId, event)
EventSource
React frontend

Camino de lectura (historial)
React
GET /notifications?cursor=&limit=
NotificationController
valida con Zod
NotificationService
.listForUser()
NotificationRepo
.findByUser(cursor)
JSON response
{ items[], nextCursor }

Establecimiento de conexión SSE
EventSource('/stream')
+ Bearer token
NotificationController
verifyJwt middleware
SseConnectionManager
.register(userId, res)
keepalive ping
cada 30 s
Al cerrar la conexión: res.on('close')SseConnectionManager.unregister(userId, res)
🔢
Secuencia de Construcción
Orden de implementación por dependencias
1
Tipos e Interfaces
notifications/types.ts lib/events.ts (tipado)
Base del sistema. Todo lo demás depende de estos contratos. Definir NotificationType, Notification, CampaignStatusChangedEvent y el mapa tipado del EventBus.
2
Migración de Base de Datos
prisma/schema.prisma
Añadir modelo Notification. Ejecutar prisma migrate dev --name add_notifications. El repositorio necesita la tabla antes de poder escribirse.
3
Capa de Infraestructura
lib/sse.ts notifications/repository.ts schemas/notification.schema.ts
Primero el SseConnectionManager (sin dependencias), luego el repositorio Prisma. Los Zod schemas se crean aquí para validar en controller y tests.
4
Lógica de Negocio + Integración
notifications/service.ts campaigns/service.ts
Implementar NotificationService con la suscripción al EventBus. Modificar CampaignService.updateStatus() para emitir el evento. Verificar el flujo completo manualmente.
5
API HTTP + Bootstrap
notifications/controller.ts src/app.ts
Crear las 4 rutas Express. Montar el router en app.ts e inicializar el NotificationService al arranque para activar la suscripción al EventBus.
6
Tests + Documentación
__tests__/notifications.test.ts
Tests unitarios (mock EventBus + mock PrismaClient) para el service. Tests de integración con supertest para el controller. Documentar el endpoint SSE en el OpenAPI/README del proyecto.