Saltar a contenido

Visión General de la Arquitectura

El middleware sigue una arquitectura basada en Sources, Destinations y Flows, orquestados por un sistema de polling y colas de trabajo.

Componentes Principales

flowchart LR
    subgraph Sources
        S1[MSSQL Source]
        S2[HTTP Source]
    end

    subgraph Services
        P[Poller]
        H[Hash Store]
        Q[Job Queue]
        M[Metrics]
    end

    subgraph Destinations
        D1[HTTP Adapter]
        D2[MySQL Adapter]
        D3[MSSQL Adapter]
    end

    S1 --> P
    S2 --> P
    P <--> H
    P --> Q
    P --> M
    Q --> D1
    Q --> D2
    Q --> D3

Flujo de Datos

sequenceDiagram
    participant Source
    participant Poller
    participant HashStore
    participant Queue
    participant Worker
    participant Destination

    loop Cada pollInterval
        Poller->>Source: getData()
        Source-->>Poller: rows[]

        loop Para cada row
            Poller->>Poller: calcular hash
            Poller->>HashStore: getHash(pk)
            HashStore-->>Poller: oldHash

            alt hash cambió
                Poller->>Queue: addJob(upsert)
            end
        end

        Poller->>HashStore: setHashes(newHashes)
    end

    Worker->>Queue: getJob()
    Queue-->>Worker: job
    Worker->>Destination: send(data)
    Destination-->>Worker: result

    alt notifyOnSuccess configurado
        Worker->>Destination: notify webhook
    end

Registries

Source Registry

Administra los adaptadores de lectura. Cada source se registra con un nombre único y puede soportar múltiples entities.

const sourceRegistry = new SourceRegistry(logger);
await sourceRegistry.registerAll([
  tangoSource,        // 'tango'
  tangoDevelopmentSource,  // 'tango-development'
  wmsSource,          // 'wms-source'
]);

Adapter Registry (Destinations)

Administra los adaptadores de escritura.

const destRegistry = new AdapterRegistry(logger);
await destRegistry.registerAll([
  wmsAdapter,         // 'wms'
  wmsDevelopmentAdapter,   // 'wms-development'
  mysqlWebAdapter,    // 'mysql-web'
  tangoAdapter,       // 'tango' (destino)
]);

Sync Flows

Un Flow define la relación entre un source y uno o más destinations:

{
  name: 'articulos',
  source: { adapter: 'tango' },
  destinations: ['wms', 'mysql-web'],
  pollInterval: 60 * 1000,    // 1 minuto
  priority: 'high',
  transform: 'miTransform',   // Opcional
}

Propiedades del Flow

Propiedad Tipo Descripción
name string Nombre único del flow (también es el entity por defecto)
source.adapter string Nombre del source adapter
source.entity string? Entity en el source (default: name)
destinations string[] Lista de destination adapters
pollInterval number Intervalo de polling en ms
priority 'high' | 'normal' | 'low' Prioridad de los jobs
transform string | string[] Transform(s) a aplicar
useLastSync boolean Usar filtro de fecha incremental

Prioridades

Prioridad Valor BullMQ Uso típico
high 1 Pedidos, artículos (cambios frecuentes)
normal 5 Clientes, proveedores
low 10 Tablas maestras (provincias, zonas)

Detección de Cambios

El sistema usa hashes MD5 para detectar cambios:

  1. Poll: Lee todos los datos del source
  2. Hash: Calcula MD5 de cada row
  3. Compara: Contra hash almacenado en Redis
  4. Encola: Solo los que cambiaron
  5. Guarda: Nuevos hashes en Redis
flowchart TD
    A[Leer datos] --> B[Calcular hash por row]
    B --> C{¿Hash existe?}
    C -->|No| D[Encolar INSERT]
    C -->|Sí| E{¿Hash cambió?}
    E -->|Sí| F[Encolar UPDATE]
    E -->|No| G[Ignorar]
    D --> H[Guardar hash]
    F --> H

Retry y Backoff

Los jobs fallidos se reintentan con backoff exponencial:

  • Intentos máximos: 5
  • Delay base: 5 segundos
  • Secuencia: 5s → 10s → 20s → 40s → 80s
flowchart LR
    A[Job falla] --> B{¿Intentos < 5?}
    B -->|Sí| C[Esperar backoff]
    C --> D[Reintentar]
    D --> E{¿Éxito?}
    E -->|No| A
    E -->|Sí| F[Completado]
    B -->|No| G[Fallo permanente]