Saltar al contenido

Patrones de Event Sourcing para Aplicaciones con Auditoría Exigente

Event sourcing en aplicaciones con auditoría exigente: almacenes de solo adición, proyecciones y snapshots que equilibran cumplimiento y rendimiento.

5 min de lectura
Diagrama que muestra un almacén de eventos de solo adición con proyecciones que se ramifican hacia diferentes modelos de lectura para auditoría y consultas

Las aplicaciones CRUD tradicionales sobreescriben el estado. Cuando alguien actualiza un registro, el valor anterior desaparece a menos que hayas añadido un registro de auditoría como algo secundario. En dominios donde cada cambio debe poder rastrearse—servicios financieros, salud, tecnología legal—esto crea una tensión fundamental entre cómo funciona la aplicación y lo que exigen los reguladores.

Event sourcing invierte el modelo de datos: en lugar de almacenar el estado actual, guardas la secuencia de eventos que lo produjo. El estado actual se convierte en una derivación, y el historial completo existe por defecto. Para aplicaciones con auditoría exigente, esto no es un patrón arquitectónico sofisticado—es la opción natural.

Eventos como fuente de verdad

En event sourcing, tu almacén de datos es un registro de solo adición de eventos de dominio. Cada evento captura qué pasó, cuándo y los detalles relevantes. El estado actual se calcula reproduciendo esos eventos.

tstypescript
// Domain events for an account management system
interface DomainEvent {
  eventId: string;
  aggregateId: string;
  eventType: string;
  timestamp: Date;
  version: number;
  data: Record<string, unknown>;
  metadata: {
    userId: string;
    correlationId: string;
    causationId: string;
    ipAddress: string;
  };
}
 
// ❌ CRUD approach: audit is an afterthought
class AccountRepository {
  async updateBalance(accountId: string, newBalance: number) {
    // Previous balance is gone forever
    await db.query(
      'UPDATE accounts SET balance = $1 WHERE id = $2',
      [newBalance, accountId]
    );
    // Audit log added later, often incomplete
    await db.query(
      'INSERT INTO audit_log (entity, action, timestamp) VALUES ($1, $2, NOW())',
      [accountId, 'balance_updated']
    );
  }
}
 
// ✅ Event sourcing: audit is inherent
class AccountEventStore {
  async append(event: DomainEvent): Promise<void> {
    await db.query(
      `INSERT INTO events (event_id, aggregate_id, event_type, 
       timestamp, version, data, metadata)
       VALUES ($1, $2, $3, $4, $5, $6, $7)`,
      [
        event.eventId,
        event.aggregateId,
        event.eventType,
        event.timestamp,
        event.version,
        JSON.stringify(event.data),
        JSON.stringify(event.metadata),
      ]
    );
  }
 
  async getEvents(aggregateId: string): Promise<DomainEvent[]> {
    const result = await db.query(
      `SELECT * FROM events 
       WHERE aggregate_id = $1 
       ORDER BY version ASC`,
      [aggregateId]
    );
    return result.rows;
  }
}

El campo metadata es crítico para la auditoría. Cada evento captura quién lo disparó, la cadena de correlación que vincula operaciones relacionadas e información contextual como direcciones IP que los auditores necesitan.

Construyendo modelos de lectura con proyecciones

Almacenar eventos resuelve el problema de auditoría pero crea un problema de consulta: no puedes consultar eficientemente "todas las cuentas con saldo superior a $10,000" reproduciendo millones de eventos cada vez. Las proyecciones resuelven esto manteniendo modelos de lectura desnormalizados.

tstypescript
// Projection that builds a read model from events
class AccountBalanceProjection {
  async handle(event: DomainEvent): Promise<void> {
    switch (event.eventType) {
      case 'AccountOpened':
        await db.query(
          `INSERT INTO account_balances 
           (account_id, balance, owner_name, opened_at)
           VALUES ($1, $2, $3, $4)`,
          [
            event.aggregateId,
            event.data.initialDeposit,
            event.data.ownerName,
            event.timestamp,
          ]
        );
        break;
 
      case 'FundsDeposited':
        await db.query(
          `UPDATE account_balances 
           SET balance = balance + $1, updated_at = $2
           WHERE account_id = $3`,
          [event.data.amount, event.timestamp, event.aggregateId]
        );
        break;
 
      case 'FundsWithdrawn':
        await db.query(
          `UPDATE account_balances 
           SET balance = balance - $1, updated_at = $2
           WHERE account_id = $3`,
          [event.data.amount, event.timestamp, event.aggregateId]
        );
        break;
    }
  }
}
 
// Separate audit-specific projection
class AuditTrailProjection {
  async handle(event: DomainEvent): Promise<void> {
    await db.query(
      `INSERT INTO audit_trail 
       (event_id, aggregate_id, event_type, actor_id,
        ip_address, correlation_id, timestamp, details)
       VALUES ($1, $2, $3, $4, $5, $6, $7, $8)`,
      [
        event.eventId,
        event.aggregateId,
        event.eventType,
        event.metadata.userId,
        event.metadata.ipAddress,
        event.metadata.correlationId,
        event.timestamp,
        JSON.stringify(event.data),
      ]
    );
  }
}

Las proyecciones son desechables. Si un auditor necesita un nuevo formato de informe, construyes una nueva proyección y reproduces los eventos a través de ella. Los datos originales nunca cambian.

Estrategias de snapshots para el rendimiento

Reproducir miles de eventos para reconstruir un agregado se vuelve lento. Los snapshots capturan periódicamente el estado calculado para que solo reproduzcas eventos posteriores al snapshot.

tstypescript
class SnapshotStore {
  private readonly SNAPSHOT_INTERVAL = 100;
 
  async getAggregate(aggregateId: string): Promise<Account> {
    const snapshot = await this.getLatestSnapshot(aggregateId);
    const startVersion = snapshot ? snapshot.version + 1 : 0;
 
    const events = await this.eventStore.getEventsAfterVersion(
      aggregateId,
      startVersion
    );
 
    let account = snapshot
      ? Account.fromSnapshot(snapshot.state)
      : Account.empty(aggregateId);
 
    for (const event of events) {
      account = account.apply(event);
    }
 
    // Create snapshot if enough events have accumulated
    if (events.length >= this.SNAPSHOT_INTERVAL) {
      await this.saveSnapshot({
        aggregateId,
        version: account.version,
        state: account.toSnapshot(),
        createdAt: new Date(),
      });
    }
 
    return account;
  }
 
  private async getLatestSnapshot(
    aggregateId: string
  ): Promise<Snapshot | null> {
    const result = await db.query(
      `SELECT * FROM snapshots 
       WHERE aggregate_id = $1 
       ORDER BY version DESC LIMIT 1`,
      [aggregateId]
    );
    return result.rows[0] || null;
  }
 
  private async saveSnapshot(snapshot: Snapshot): Promise<void> {
    await db.query(
      `INSERT INTO snapshots 
       (aggregate_id, version, state, created_at)
       VALUES ($1, $2, $3, $4)`,
      [
        snapshot.aggregateId,
        snapshot.version,
        JSON.stringify(snapshot.state),
        snapshot.createdAt,
      ]
    );
  }
}

Los snapshots son una optimización, no datos originales. Siempre puedes eliminar todos los snapshots y reconstruir desde los eventos: el sistema sigue siendo correcto, solo temporalmente más lento.

Manejando la evolución del esquema de eventos

Los eventos son inmutables, pero tu comprensión del dominio evoluciona. Necesitarás lidiar con eventos escritos con un esquema anterior al que tu código actual espera.

tstypescript
// Event upcasting: transform old event shapes to new ones
class EventUpcaster {
  private upcasters = new Map<string, UpcastFunction[]>();
 
  register(eventType: string, fromVersion: number, upcaster: UpcastFunction) {
    const key = `${eventType}:${fromVersion}`;
    const chain = this.upcasters.get(key) || [];
    chain.push(upcaster);
    this.upcasters.set(key, chain);
  }
 
  upcast(event: StoredEvent): DomainEvent {
    let current = event;
 
    while (current.schemaVersion < this.getCurrentVersion(current.eventType)) {
      const key = `${current.eventType}:${current.schemaVersion}`;
      const upcaster = this.upcasters.get(key);
      if (!upcaster) {
        throw new Error(
          `No upcaster for ${current.eventType} v${current.schemaVersion}`
        );
      }
      current = upcaster[0](current);
    }
 
    return current as DomainEvent;
  }
}
 
// Example: FundsDeposited event evolved over time
// v1: { amount: number }
// v2: { amount: number, currency: string }
// v3: { amount: number, currency: string, channel: string }
 
const upcaster = new EventUpcaster();
 
upcaster.register('FundsDeposited', 1, (event) => ({
  ...event,
  data: { ...event.data, currency: 'USD' },
  schemaVersion: 2,
}));
 
upcaster.register('FundsDeposited', 2, (event) => ({
  ...event,
  data: { ...event.data, channel: 'unknown' },
  schemaVersion: 3,
}));

Consultas de cumplimiento: el beneficio de la auditoría

Con event sourcing, las consultas de cumplimiento que requerirían joins complejos de registros de auditoría en un sistema CRUD se convierten en consultas directas al flujo de eventos.

tstypescript
class ComplianceQueryService {
  // "Show me every change to this account, by whom, and when"
  async getAccountHistory(accountId: string): Promise<AuditEntry[]> {
    const events = await this.eventStore.getEvents(accountId);
    return events.map((event) => ({
      timestamp: event.timestamp,
      action: event.eventType,
      actor: event.metadata.userId,
      ipAddress: event.metadata.ipAddress,
      details: event.data,
      correlationId: event.metadata.correlationId,
    }));
  }
 
  // "What did the account look like at this specific point in time?"
  async getStateAtTime(
    accountId: string,
    pointInTime: Date
  ): Promise<Account> {
    const events = await this.eventStore.getEventsUntil(
      accountId,
      pointInTime
    );
    let account = Account.empty(accountId);
    for (const event of events) {
      account = account.apply(event);
    }
    return account;
  }
 
  // "Who accessed accounts with balances over $50k last quarter?"
  async getHighValueAccessPatterns(
    threshold: number,
    startDate: Date,
    endDate: Date
  ): Promise<AccessPattern[]> {
    // This query runs against a projection built
    // specifically for compliance reporting
    return db.query(
      `SELECT actor_id, account_id, event_type, timestamp
       FROM compliance_access_log
       WHERE balance_at_time > $1
       AND timestamp BETWEEN $2 AND $3
       ORDER BY timestamp DESC`,
      [threshold, startDate, endDate]
    );
  }
}

La consulta de "estado en un punto en el tiempo" es la función destacada para la auditoría. En un sistema CRUD, reconstruir el estado histórico requiere una correlación extensa de registros de auditoría. En event sourcing, es reproducir eventos hasta un timestamp.

Puntos clave

Event sourcing convierte los rastros de auditoría en una consecuencia natural del modelo de datos en lugar de algo añadido a contrarreloj: cada cambio de estado es un evento inmutable que captura no solo qué cambió, sino quién lo disparó, cuándo y por qué. Las proyecciones resuelven el problema de consultas manteniendo modelos de lectura desechables: construye nuevas proyecciones cuando los auditores necesiten una vista diferente, reproduce el historial de eventos a través de ellas y nunca te preocupes por la pérdida de datos porque los eventos originales son inmutables. Los snapshots son una optimización de rendimiento que te permite omitir la reproducción de eventos antiguos: son resúmenes desechables, no datos originales, y solo debes introducirlos cuando el rendimiento de la reproducción se vuelva realmente un problema, no de forma preventiva. La evolución del esquema mediante upcasting permite que tu modelo de dominio crezca sin reescribir los eventos almacenados: transforma formas antiguas de eventos a esquemas actuales en tiempo de lectura, manteniendo la compatibilidad hacia atrás con todo tu historial de eventos.

Wilfredo Rujel

Wilfredo Rujel

Ingeniero de Software Full Stack

Compartir esta publicaciónX