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:
- Poll: Lee todos los datos del source
- Hash: Calcula MD5 de cada row
- Compara: Contra hash almacenado en Redis
- Encola: Solo los que cambiaron
- 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]