Modelo de Dominio Event-First
Filosofía
El diseño event-first significa que los eventos de dominio son ciudadanos de primera clase. Cada acción significativa en el sistema produce un evento. Las proyecciones se derivan de eventos. El estado se reconstruye desde eventos. Esta arquitectura facilita:
- Auditoría completa: cada cambio queda registrado.
- Reactivity: componentes reaccionan a cambios en tiempo real.
- Desacoplamiento: productores no conocen a consumidores.
- Replayability: el estado puede reconstruirse desde el historial de eventos.
- Extensibilidad: nuevos consumidores se agregan sin modificar productores.
Eventos de Dominio Descubiertos
Contactos
| Evento | Descripción | Datos esenciales |
|---|---|---|
ContactCreated | Nuevo contacto creado | contact_id, tenant_id, attributes |
ContactUpdated | Contacto modificado | contact_id, changed_fields, previous_values |
ContactMerged | Dos contactos fusionados | primary_contact_id, merged_contact_id, merged_attributes |
Mensajería
| Evento | Descripción | Datos esenciales |
|---|---|---|
MessageReceived | Mensaje entrante de un canal | message_id, conversation_id, channel, sender, content, timestamp |
MessageSent | Mensaje enviado por un agente | message_id, conversation_id, agent_id, content, channel |
MessageDelivered | Confirmación de entrega del canal | message_id, channel, delivered_at |
MessageRead | El destinatario leyó el mensaje | message_id, read_at |
Conversaciones
| Evento | Descripción | Datos esenciales |
|---|---|---|
ConversationOpened | Nueva conversación iniciada | conversation_id, channel, contact_id, initial_message |
ConversationAssigned | Conversación asignada a un agente | conversation_id, agent_id, assigned_by, method |
ConversationClosed | Conversación finalizada | conversation_id, closed_by, disposition, resolution_time |
ConversationTransferred | Conversación transferida a otro agente o equipo | conversation_id, from_agent_id, to_agent_id, reason |
Agentes
| Evento | Descripción | Datos esenciales |
|---|---|---|
AgentAvailable | Agente marcado como disponible | agent_id, timestamp |
AgentUnavailable | Agente marcado como no disponible | agent_id, reason, timestamp |
AgentPaused | Agente en pausa (break) | agent_id, pause_type, timestamp |
AgentResumed | Agente reanudó después de pausa | agent_id, resume_reason, duration |
Campañas
| Evento | Descripción | Datos esenciales |
|---|---|---|
CampaignStarted | Campañas iniciada | campaign_id, started_by, target_count |
CampaignPaused | Campañas pausada | campaign_id, paused_by, reason |
CampaignCompleted | Campañas finalizada | campaign_id, stats (contacted, answered, converted) |
Telefonía
| Evento | Descripción | Datos esenciales |
|---|---|---|
CallInitiated | Llamada iniciada | call_id, direction, from, to, tenant_id |
CallAnswered | Llamada contestada | call_id, answered_by, answer_time |
CallEnded | Llamada finalizada | call_id, duration, disposition, hangup_cause |
CallRecordingStarted | Grabación iniciada | call_id, recording_id, started_by |
Tickets
| Evento | Descripción | Datos esenciales |
|---|---|---|
TicketCreated | Ticket de soporte creado | ticket_id, conversation_id, priority, category |
TicketUpdated | Ticket actualizado | ticket_id, changed_fields, updated_by |
TicketResolved | Ticket resuelto | ticket_id, resolution, resolved_by, resolution_time |
Automatización
| Evento | Descripción | Datos esenciales |
|---|---|---|
AutomationTriggered | Regla de automatización ejecutada | automation_id, trigger_event, conditions_matched, actions_executed |
Disposición
| Evento | Descripción | Datos esenciales |
|---|---|---|
DispositionRecorded | Disposición registrada al cerrar conversación | conversation_id, agent_id, disposition_code, notes |
Commands
Los commands representan intenciones de cambio en el estado del sistema. Cada command produce uno o más eventos.
Mensajería
| Command | Descripción | Eventos producidos |
|---|---|---|
SendMessage | Enviar mensaje a través de un canal | MessageSent, (async) MessageDelivered |
TransferConversation | Transferir conversación a otro agente o equipo | ConversationTransferred, ConversationAssigned |
Gestión de agentes
| Command | Descripción | Eventos producidos |
|---|---|---|
AssignAgent | Asignar conversación a un agente | ConversationAssigned |
PauseAgent | Poner agente en pausa | AgentPaused |
Campañas
| Command | Descripción | Eventos producidos |
|---|---|---|
StartCampaign | Iniciar campaña | CampaignStarted |
PauseCampaign | Pausar campaña | CampaignPaused |
StopCampaign | Detener campaña | CampaignCompleted |
Telefonía
| Command | Descripción | Eventos producidos |
|---|---|---|
OriginateCall | Iniciar llamada saliente | CallInitiated |
HangupCall | Finalizar llamada | CallEnded |
BridgeCall | Unir dos canales (conferencia) | (evento de bridge) |
Contactos
| Command | Descripción | Eventos producidos |
|---|---|---|
CreateContact | Crear nuevo contacto | ContactCreated |
UpdateContact | Actualizar contacto existente | ContactUpdated |
MergeContacts | Fusionar dos contactos | ContactMerged |
Aggregates (Raíces de Agregado)
Conversation (raíz de agregado para mensajería)
El agregado Conversation es el centro del sistema de mensajería:
Estado:
id: identificador único.tenant_id: multi-tenancy.channel: canal de comunicación.contact_id: contacto asociado.status: opened, assigned, pending, resolved, closed.assigned_agent_id: agente actual (nullable).priority: prioridad (normal, high, urgent).tags: etiquetas para clasificación.sla: SLA aplicable.created_at,updated_at,closed_at.
Invariants:
- Una conversación solo puede estar asignada a un agente a la vez.
- Una conversación cerrada no puede reabrirse (se crea una nueva).
- El SLA se evalúa contra el timestamp del último mensaje del cliente.
- Solo agentes disponibles pueden recibir asignaciones.
Commands que acepta:
AssignAgent: asigna un agente si la conversación está abierta.Transfer: transfiere a otro agente.Close: cierra la conversación con disposición.AddMessage: agrega un mensaje a la conversación.
Campaign (raíz de agregado para outbound)
Estado:
id: identificador único.tenant_id: multi-tenancy.name: nombre de la campaña.status: draft, active, paused, completed, cancelled.channel: canal de la campaña.pacing_config: configuración de velocidad.schedule: horario de ejecución.target_count: número de contactos objetivo.stats: estadísticas en tiempo real.
Invariants:
- Una campaña solo puede estar activa una vez.
- El pacing no puede exceder los límites configurados.
- Las campañas respetan horarios de quieter hours.
- DNC list se verifica antes de cada intento.
Agent (raíz de agregado para workforce)
Estado:
id: identificador único.tenant_id: multi-tenancy.status: offline, available, on_call, paused, wrap_up.skills: habilidades para routing.current_conversations: conversaciones asignadas.max_conversations: límite de conversaciones simultáneas.channels: canales que puede atender.pause_reason: razón de pausa actual.
Invariants:
- Un agente no puede recibir más conversaciones que su límite.
- Un agente en pausa no puede recibir nuevas asignaciones.
- Un agente offline no puede ser asignado.
- Las habilidades se usan para skills-based routing.
Contact (raíz de agregado para CRM)
Estado:
id: identificador único.tenant_id: multi-tenancy.external_ids: mapeo de IDs de canales (phone, email, whatsapp, etc.).attributes: atributos personalizados (JSONB).company_id: empresa asociada (nullable).interactions_count: contador de interacciones.last_interaction_at: timestamp de última interacción.created_at,updated_at.
Invariants:
- Los external_ids deben ser únicos por canal dentro del tenant.
- El merge de contactos preserva el contacto primario.
- Los atributos personalizados se validan contra el schema del tenant.
Tenant (raíz de agregado para multi-tenancy)
Estado:
id: identificador único.name: nombre del tenant.status: active, suspended, trial, cancelled.plan: plan de suscripción.config: configuración del tenant (JSONB).limits: límites de recursos (canales, agentes, storage).created_at,updated_at.
Invariants:
- Un tenant suspendido no puede crear nuevos recursos.
- Los límites se verifican antes de crear recursos.
- La configuración del tenant se propaga a todos los contextos.
Proyecciones
Las proyecciones son vistas materializadas derivadas de eventos. Se actualizan cuando se procesan eventos.
ConversationList (para vista de inbox)
Fuente: ConversationOpened, ConversationAssigned, ConversationClosed, MessageReceived, MessageSent.
Campos: id, contact_name, channel, status, assigned_agent, last_message_preview, last_message_at, unread_count, priority, tags.
Uso: inbox de agentes, lista de conversaciones activas.
AgentStatus (para monitoreo en tiempo real)
Fuente: AgentAvailable, AgentUnavailable, AgentPaused, AgentResumed, ConversationAssigned, ConversationClosed.
Campos: agent_id, name, status, current_conversations, max_conversations, last_status_change, pause_duration_today, conversations_handled_today.
Uso: dashboard de supervisores, real-time monitoring.
CampaignStats (para dashboard)
Fuente: CampaignStarted, CampaignPaused, CampaignCompleted, plus call/message events.
Campos: campaign_id, contacted, answered, converted, pending, conversion_rate, avg_duration, cost_per_contact.
Uso: dashboards de campañas, reporting en tiempo real.
ReportData (para analytics)
Fuente: todos los eventos relevantes, procesados en batch.
Campos: various metrics aggregated by time period, agent, campaign, channel, etc.
Uso: reportes históricos, análisis de tendencias, exportaciones.
Procesamiento
Inline (síncrono)
Operaciones que se procesan inmediatamente dentro de la misma transacción:
| Evento | Procesamiento |
|---|---|
MessageReceived | Actualización de ConversationList, incremento de unread_count |
ConversationAssigned | Actualización de AgentStatus, notificación al agente |
AgentAvailable | Actualización de routing availability |
AgentPaused | Remover del routing pool |
Garantía: consistencia fuerte dentro de la transacción.
Asíncrono
Operaciones que se procesan después, sin bloquear el request:
| Evento | Procesamiento |
|---|---|
MessageReceived | Notificación push al agente, webhook a integración externa |
ConversationClosed | Generación de reporte, envío de email de encuesta, actualización de CRM |
CallEnded | Procesamiento de grabación, generación de transcript, actualización de métricas |
CampaignStarted | Inicio de discado, scheduling de llamadas |
Garantía: eventual consistency. Los eventos se procesan en orden pero puede haber delay.
Idempotencia
- Message deduplication: cada mensaje tiene un
external_message_idúnico del canal. Se verifica antes de procesar para evitar duplicados. - Command idempotency: cada command tiene un
idempotency_keyúnico. Se usa para prevenir procesamiento duplicado. - Event deduplication: los eventos tienen un
event_idúnico. Los consumidores verifican antes de procesar.
Ordering
- Per-conversation: los eventos dentro de una conversación se procesan en orden estricto. Se usa una cola por conversation_id.
- Cross-conversation: entre conversaciones, se garantiza eventual consistency. No hay orden garantizado entre conversaciones diferentes.
- Implementation: cola de mensajes con partitioning por conversation_id (Redis Streams o PostgreSQL con advisory locks).
Infraestructura de Eventos
No Kafka inicialmente
Kafka es una solución potente pero introduce complejidad operativa significativa. Para la fase inicial:
- No se justifica el overhead operativo de Kafka.
- Alternativas más simples cubren los requisitos iniciales.
- Migración futura es posible si el volumen lo requiere.
PostgreSQL-backed event store
Almacenamiento de eventos en PostgreSQL:
CREATE TABLE domain_events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
tenant_id UUID NOT NULL,
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id UUID NOT NULL,
event_type VARCHAR(100) NOT NULL,
event_data JSONB NOT NULL,
metadata JSONB,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
processed_at TIMESTAMPTZ
);
CREATE INDEX idx_events_aggregate ON domain_events(aggregate_type, aggregate_id);
CREATE INDEX idx_events_tenant ON domain_events(tenant_id, created_at);
CREATE INDEX idx_events_type ON domain_events(event_type, created_at);Ventajas:
- Transaccionalidad con el resto de la base de datos.
- Consultas SQL estándar para reporting.
- Sin infraestructura adicional.
Desventajas:
- No escala a alto volumen como Kafka.
- Polling para consumidores (no push nativo).
- Latencia de procesamiento depende de la frecuencia de polling.
Redis Streams como alternativa
Para eventos que requieren menor latencia:
XADD events * tenant_id "t1" event_type "MessageReceived" data '{"message_id":"m1"}'
XREAD COUNT 10 BLOCK 5000 STREAMS events 0Ventajas:
- Push nativo (XREAD BLOCK).
- Consumer groups para procesamiento distribuido.
- Retención configurable de eventos.
- Menor latencia que polling.
Desventajas:
- No transaccional con PostgreSQL.
- Persistencia configurable (puede perder datos si no se configura correctamente).
Outbox Pattern
Para publicación confiable de eventos:
- El dominio escribe el evento en una tabla
outboxdentro de la misma transacción que el cambio de estado. - Un proceso poller lee la tabla
outboxy publica los eventos al canal de distribución (Redis Streams, o directamente a consumidores). - Los eventos publicados se marcan como
publisheden la tablaoutbox.
CREATE TABLE outbox (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_type VARCHAR(100) NOT NULL,
aggregate_id UUID NOT NULL,
event_type VARCHAR(100) NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
published_at TIMESTAMPTZ
);Garantía: los eventos se publican al menos una vez (at-least-once delivery). Los consumidores deben ser idempotentes.
Process Managers para workflows multi-paso
Un process manager orquesta flujos que requieren múltiples pasos:
Ejemplo: cierre de conversación:
ConversationClosedevent received.- Process manager ejecuta: generar reporte → enviar encuesta → actualizar CRM → liberar agente.
- Si falla un paso, se reintenta con backoff exponencial.
- Si falla definitivamente, se notifica para intervención manual.
Implementación: process managers se ejecutan como async workers. Cada process manager tiene su lógica de retry y compensación.