Skip to content

Plataforma técnica · Orchestrator

Domain Events

Orchestrator Common Events

El bus de eventos de dominio es el mecanismo que usa el Orchestrator para que cambios en un servicio disparen reacciones en otros sin acoplarlos directamente. Es in-process, sincrónico para listeners de negocio, fire-and-forget para sinks de observabilidad, y vive en domain/common/events/.

flowchart LR
  subgraph EMI["Emisores"]
    CAP["CapitalService"]
    PAY["PayrollService"]
    CICLO["CicloContableService"]
    OTROS["..."]
  end

  subgraph TX["Outbox (TX)"]
    OB["queue + drain<br/>(post-COMMIT)"]
  end

  subgraph BUS["DomainEventBus<br/>(singleton)"]
    LST["listeners<br/>Map&lt;event, Set&lt;handler&gt;&gt;"]
    SNK["sinks<br/>Set&lt;EventSink&gt;"]
    HIS["history (max 1000)"]
  end

  subgraph LIS["Listeners (negocio)"]
    CACHE["cache.listeners<br/>invalida apiCache"]
    FIN["financieros.listeners<br/>match conciliación"]
    L2["..."]
  end

  subgraph SNKS["Sinks (observabilidad)"]
    MGO["MongoEventSink<br/>nostromo_cortex.domain_events"]
  end

  CAP & PAY & CICLO & OTROS --> OB --> BUS
  BUS --> LST --> LIS
  BUS --> SNK --> SNKS
type EventContext = {
tenantDb: string;
userId: string;
requestId?: string;
};
type EventPayload<T> = T & EventContext; // payload de negocio + ctx
interface DomainEventMap {
'capital:movimiento': { capitalId, tipoMovimiento, monto };
'payroll:approved': { payrollId, employeeId, period, totalLiquid };
// ... ver catálogo completo
}
type DomainEventName = keyof DomainEventMap;
type EventHandler<T> = (payload: EventPayload<DomainEventMap[T]>) => void | Promise<void>;
type EventSink = (envelope: EventEnvelope) => void | Promise<void>;

El DomainEventMap es la fuente de verdad del contrato — cada evento tiene payload tipado. Los listeners reciben payload enriquecido con tenantDb/userId/requestId; los sinks reciben un EventEnvelope con event, payload, ts.

MétodoPara
DomainEventBus.getInstance()Singleton accessor.
bus.on(event, handler)Suscribir un listener. Retorna unsubscribe.
bus.once(event, handler)Suscribir y auto-remove tras primer disparo.
bus.off(event, handler)Desuscribir.
bus.addSink(sink)Suscribir un sink (recibe TODOS los eventos).
bus.removeSink(sink)Desuscribir un sink.
bus.emit(event, payload)Emitir. Casi nadie llama esto directamenteBaseService.flushOutbox lo hace post-COMMIT.
bus.getHistory(event?)Últimos N eventos (max 1000 en memoria, ring buffer).
DomainEventBus.resetInstance()Solo para tests.
AspectoListenersSinks
GranularidadPor evento específico (on('payroll:approved', ...)).Reciben TODOS los eventos.
NaturalezaLógica de negocio reactiva.Observabilidad (audit, métricas, log).
Asyncawait-eados; el bus espera a todos.Fire-and-forget; el bus no espera.
ErroresLogueados, no se propagan al emisor.Logueados, no se propagan al emisor.

emit() espera a los listeners (Promise.all) antes de retornar; los sinks corren en paralelo en background.

El DomainEventMap declara 22 eventos. Resumen por dominio:

EventoPayloadEmisor
capital:movimiento{ capitalId, tipoMovimiento: APORTE|RETIRO|AJUSTE, monto }CapitalService.create
EventoPayloadEmisor
payroll:approved{ payrollId, employeeId, period: { year, month }, totalLiquid }PayrollService.approve
honorario:contabilizado{ honorarioId, rutPrestador, monto, retencion, periodo }HonorariosService.contabilizar
honorario:reversado{ honorarioId, rutPrestador, prevEstado }HonorariosService.reversar
EventoPayloadEmisor
compra:contabilizada{ compraId, monto, formaPago }OperacionesSiiService.compra
venta:contabilizada{ ventaId, documentoTipo: BOLETA|FACTURA, rutCliente, monto, formaPago }OperacionesSiiService.venta
EventoPayload
compra_inv:contabilizada{ movimientoCostoId, compraDetalleId, categoriaId, monto }
costo:egreso_calculado{ movimientoCostoId, categoriaId, periodo, monto, factorProrrateo }
existencia:castigada{ movimientoCostoId, categoriaId, monto, motivo }
EventoPayload
depreciacion:calculada{ activoFijoId, periodo, monto, vidaUtilRestanteMeses }
EventoPayload
conciliacion:matched{ movimientoBancarioId, documentoTipo, documentoId, monto, formaResolucion: AUTOMATICA|MANUAL|PARCIAL }
movimiento:vinculado{ movimientoBancarioId, categoriaMovimientoCodigo, documentoVinculadoId }
conciliacion:desconciliada{ movimientoBancarioId, prevDocumentoId }
EventoPayload
ciclo:cierre{ anio, mes: number|null, pasosCompletados: string[] }
EventoPayload
f29:generada{ declaracionId, periodo, rutContribuyente }
f29:status_cambiado{ declaracionId, periodo, estado, prevEstado }
dj:1879_generada{ declaracionId, anio, totalReceptores, totalRetencionHonorarios }
dj:1879_status_cambiado{ declaracionId, anio, estado, prevEstado }
dj:1887_generada{ declaracionId, anio, totalTrabajadores, totalRentaActualizada, totalIUSC }
dj:1887_status_cambiado{ declaracionId, anio, estado, prevEstado }
EventoPayload
user:login{ email, method: password|totp }

No llames bus.emit() directamente en código de servicio. Usá el outbox de BaseService.withTransaction:

Patrón outbox
return this.withTransaction(ctx, async (client, outbox) => {
const entry = await Repository.create(client, dto);
outbox.queue('capital:movimiento', {
capitalId: entry.id,
tipoMovimiento: 'APORTE',
monto: entry.monto,
});
return this.success(entry);
});

Razón: si la TX rollback, el outbox se descarta y el evento nunca se emite. Llamar bus.emit() antes del COMMIT puede notificar listeners de un cambio que jamás persistió.

Detalle del flujo y enriquecimiento del payload en BaseService › Outbox Transaccional.

Los listeners de negocio se registran al bootstrap del servidor (en events/bootstrap.ts), no por servicio:

Listener de cache invalidation
// orchestrator/src/domain/cache/listeners.ts
import { DomainEventBus } from '@/domain/common/events/DomainEventBus';
import { apiCache } from '@/lib/cache';
export function registerCacheListeners() {
const bus = DomainEventBus.getInstance();
bus.on('capital:movimiento', (payload) => {
apiCache.delByPattern(`capital:${payload.tenantDb}:*`);
});
// ... más listeners
}

Los listeners son síncronos respecto a la emisión — el bus espera a que todos retornen antes de seguir. Si necesitas que un listener no bloquee, hacé el trabajo dentro de setImmediate o queueá a una cola externa.

Sink default que persiste todos los eventos en nostromo_cortex.domain_events para análisis posterior (visualizable desde la extensión Cortex en VS Code).

VariableDefaultNotas
EVENTS_MONGO_ENABLEDhabilitado0 desactiva el sink al bootstrap.
EVENTS_MONGO_URImongodb://localhost:27017
EVENTS_MONGO_DBnostromo_cortex
EVENTS_MONGO_COLLECTIONdomain_events
EVENTS_MONGO_TTL_DAYS30TTL del índice sobre ts — auto-purge.
domain_events doc shape
{
ts: Date, // momento de emisión
event: string, // e.g. "capital:movimiento"
source: "orchestrator",
tenantDb: string,
userId: string,
requestId?: string,
payload: { /* sin tenantDb/userId/requestId */ },
host: string, // os.hostname()
pid: number // process.pid
}

stripContext elimina los campos de EventContext del payload antes de persistir — ya están como columnas top-level.

  • { ts: 1 } con expireAfterSeconds = ttlSeconds (TTL natural de Mongo).
  • { tenantDb: 1, event: 1, ts: -1 } para consultas comunes desde Cortex.

Si Mongo no responde al bootstrap (serverSelectionTimeoutMS: 2000):

  1. Log una sola advertencia ([MongoEventSink] no se pudo conectar a ..., sink quedará en no-op).
  2. collection = null — el sink se vuelve un no-op silencioso.
  3. Cada insertOne posterior chequea if (!collection) return.

Si Mongo está OK al inicio pero falla durante runtime, insertOne captura el error y loguea una sola vez (warned = true), luego silencia para no inundar stdout. El emisor nunca ve este error.

initializeEventBus() en events/bootstrap.ts debe llamarse una vez al arrancar el servidor:

orchestrator/src/server.ts
import { initializeEventBus } from '@/domain/common/events/bootstrap';
await initializeEventBus();

Hace tres cosas:

  1. initializeDomainListeners() — registra los listeners de dominio (registerCacheListeners, registerFinancierosListeners, …).
  2. Si EVENTS_MONGO_ENABLED !== '0': crea MongoEventSink y lo añade al bus.
  3. Loguea [EventBus] initialized.

shutdownEventBus() cierra la conexión Mongo del sink. Conviene llamarlo en el handler beforeExit/SIGTERM.

Para tests, llamar DomainEventBus.resetInstance() en beforeEach evita que listeners de un test contaminen el siguiente.

Para verificar que un servicio emitió un evento esperado, suscribir un mock listener:

Test pattern
it('emite capital:movimiento tras crear', async () => {
const handler = jest.fn();
DomainEventBus.getInstance().on('capital:movimiento', handler);
await capitalService.create(ctx, dto);
expect(handler).toHaveBeenCalledWith(
expect.objectContaining({ tipoMovimiento: 'APORTE' }),
);
});