Saltar a contenido

Destinations (Adaptadores de Escritura)

Los destinations son adaptadores que escriben datos a sistemas externos. Hay tres tipos implementados.

Tipos de Destinations

classDiagram
    class Adapter {
        <<interface>>
        +name: string
        +connect(): Promise~void~
        +send(entity, pk, data): Promise~AdapterResult~
        +delete(entity, pk): Promise~AdapterResult~
        +supports(entity): boolean
        +close(): Promise~void~
    }

    class BaseAdapter {
        #mapData(data, fields, transform): RowData
        #success(): AdapterResult
        #error(msg): AdapterResult
    }

    class HttpAdapter {
        -config: HttpConfig
        +send()
        +delete()
    }

    class MysqlAdapter {
        -pool: Pool
        +send()
        +delete()
    }

    class MssqlDestAdapter {
        -pool: ConnectionPool
        +send()
        -notify()
    }

    Adapter <|.. BaseAdapter
    BaseAdapter <|-- HttpAdapter
    BaseAdapter <|-- MysqlAdapter
    BaseAdapter <|-- MssqlDestAdapter

HTTP Adapter

Escribe datos a APIs REST (WMS).

Configuración

// src/config/destinations/wms.ts
export const wmsAdapter: HttpAdapterConfig = {
  name: 'wms',
  type: 'http',
  config: {
    baseUrl: process.env.WMS_API_URL || 'http://localhost:8888',
    timeout: 30000,
    headers: {
      'Content-Type': 'application/json',
      'X-API-KEY': process.env.WMS_API_KEY || '',
    },
  },
  entities: {
    articulos: {
      endpoint: 'articulos',
      keyField: 'codigo',
      fields: {
        codigo: 'CODIGO',
        nombre: 'DESCRIPCION',
        // origen → destino
      },
    },
    // ... más entities
  },
};

Entity Config HTTP

Propiedad Tipo Descripción
endpoint string Path del endpoint
keyField string Campo usado como ID en la URL
fields Record Mapeo de campos (destino → origen)
transform string? Transform a aplicar

Lógica de Upsert

El adapter HTTP intenta POST primero, si recibe 409 (Conflict) hace PATCH:

flowchart TD
    A[send] --> B[POST /endpoint]
    B --> C{Status?}
    C -->|200/201| D[✅ Success]
    C -->|409| E[PATCH /endpoint/id]
    E --> F{Status?}
    F -->|200| D
    F -->|Error| G[❌ Error]
    C -->|Otro| G

MySQL Adapter

Escribe datos a MySQL/MariaDB. Soporta dos modos:

  • Una entity → una tabla (caso clásico): un INSERT/UPSERT/REPLACE por cada fila.
  • Una entity → múltiples tablas (multi-op, transaccional): la misma fila origen genera varias escrituras coordinadas, en una transacción, con resolución de IDs autogenerados vía subselects.

Configuración (caso clásico)

// src/config/destinations/mysql-web.ts
export const mysqlWebAdapter: MysqlAdapterConfig = {
  name: 'mysql-web',
  type: 'mysql',
  config: {
    host: process.env.WEB_MYSQL_HOST || 'localhost',
    port: parseInt(process.env.WEB_MYSQL_PORT || '3306'),
    database: process.env.WEB_MYSQL_DATABASE || 'tienda_web',
    user: process.env.WEB_MYSQL_USER || 'root',
    password: process.env.WEB_MYSQL_PASSWORD || '',
  },
  entities: {
    articulos: {
      table: 'WEBVIEW',
      mode: 'upsert',
      upsertKey: ['CODIGO'],
      fields: {
        CODIGO: 'CODIGO',
        DESCRIPCION: 'DESCRIPCION',
        PRECIO: 'PRECIO',
        // ... más campos
      },
    },
  },
};

Entity Config MySQL

Propiedad Tipo Descripción
table string? Nombre de la tabla. Opcional si se usa operations.
mode 'insert' | 'upsert' | 'replace' Modo de escritura. Opcional si se usa operations.
upsertKey string[]? Columnas que disparan ON DUPLICATE KEY UPDATE
fields Record Mapeo de campos para una sola tabla
transform string | string[]? Transform(s) a aplicar antes del mapeo
operations MysqlOperation[]? Operaciones múltiples (ver abajo). Si está, ignora table/mode/fields

Escape de Columnas

El adapter escapa automáticamente nombres de columnas con caracteres especiales:

// Columnas como "ESPECIFICACION/TIPO" o "OEM 1"
// se escapan automáticamente con backticks
`ESPECIFICACION/TIPO`  `\`ESPECIFICACION/TIPO\``

Operaciones Múltiples (multi-op transaccional)

Cuando una sola fila del source debe materializarse en varias tablas del destino MySQL — por ejemplo: dar de alta el cliente, el vendedor y la relación entre ambos a partir de un único registro — se usa operations.

Garantías:

  • Las operaciones corren en orden, en una transacción (BEGIN → ops → COMMIT).
  • Si cualquier operación falla, se hace ROLLBACK automático.
  • Las operaciones posteriores ven las escrituras de las anteriores (read-your-writes en la misma conexión), por lo que se puede resolver IDs autogenerados con SELECT.
  • El transform se aplica una sola vez al row antes de la primera op.

Caso real: clientes → MySQL Web

Una fila de WEBCLIENTE (Tango) trae datos del cliente y del vendedor mezclados. En MySQL Web hay que materializar tres registros: el cliente en REG_CUENTA (subsistema 02), el vendedor en REG_CUENTA (subsistema 22) y la relación entre ambos en REG_CUENTA_RELACIONES, donde la FK es el id_regcuenta autoincremental — que solo existe después del upsert.

// src/config/destinations/mysql-web.ts
entities: {
  clientes: {
    transform: 'clientesMysql',
    operations: [
      // 1) Cliente
      {
        table: 'REG_CUENTA',
        mode: 'upsert',
        upsertKey: ['subsistema_id', 'nro_cuenta'],
        fields: {
          subsistema_id: 'subsistema_id',     // '02' (lo setea el transform)
          nombre_cuenta: 'nombre_cuenta',
          nro_cuenta: 'nro_cuenta',
          // ... resto de campos del cliente
        },
      },
      // 2) Vendedor
      {
        table: 'REG_CUENTA',
        mode: 'upsert',
        upsertKey: ['subsistema_id', 'nro_cuenta'],
        fields: {
          subsistema_id: () => '22',
          nro_cuenta: '_vend_cod',
          nombre_cuenta: '_vend_nombre',
        },
      },
      // 3) Relación cliente ↔ vendedor (resuelve los IDs autogenerados)
      {
        table: 'REG_CUENTA_RELACIONES',
        mode: 'upsert',
        upsertKey: [
          'regcuenta_id',
          'regcuentarelacionada_id',
          'subsistemarelacionada_id',
        ],
        fields: {
          regcuenta_id: {
            select: {
              table: 'REG_CUENTA',
              column: 'id_regcuenta',
              where: { subsistema_id: () => '02', nro_cuenta: 'nro_cuenta' },
            },
          },
          regcuentarelacionada_id: {
            select: {
              table: 'REG_CUENTA',
              column: 'id_regcuenta',
              where: { subsistema_id: () => '22', nro_cuenta: '_vend_cod' },
            },
          },
          subsistemarelacionada_id: () => '22',
        },
      },
    ],
  },
}

Tipo MysqlFieldValue — qué puede ir como valor de un field:

Forma Significado
'columna' Toma row['columna'] del row transformado
(row) => valor Función que computa el valor
{ select: { table, column, where } } Subselect que se ejecuta dentro de la transacción y devuelve el valor

Subselect (select):

  • where es un mapa columna → valor; cada valor sigue las mismas reglas que un field (string lookup o función).
  • Si el SELECT devuelve 0 filas, el adapter lanza Lookup sin resultados: ... y la transacción se hace rollback.
  • Como la conexión es la misma que ejecutó las ops anteriores, ve las escrituras pendientes (read-your-writes).

Diagrama del flujo multi-op:

sequenceDiagram
    participant W as Worker
    participant P as Pool
    participant C as PoolConnection
    participant MySQL

    W->>P: getConnection()
    P-->>W: conn (reservada)
    W->>C: BEGIN
    W->>C: INSERT REG_CUENTA (cliente, sub 02) ON DUP KEY UPDATE
    C->>MySQL: ✓
    W->>C: INSERT REG_CUENTA (vendedor, sub 22) ON DUP KEY UPDATE
    C->>MySQL: ✓
    W->>C: SELECT id_regcuenta WHERE sub=02, nro=...
    C-->>W: 100
    W->>C: SELECT id_regcuenta WHERE sub=22, nro=...
    C-->>W: 200
    W->>C: INSERT REG_CUENTA_RELACIONES (100, 200, '22')
    C->>MySQL: ✓
    W->>C: COMMIT
    W->>P: release()

Por qué getConnection() y no pool.execute()

Las transacciones en MySQL son por conexión. pool.execute() toma una conexión efímera del pool por cada query — no se puede mantener una transacción abierta entre llamadas. Por eso el adapter reserva una PoolConnection específica al inicio del multi-op y la devuelve con release() al final.

Qué tiene que existir en la base

Para que el ON DUPLICATE KEY UPDATE funcione correctamente hace falta:

  • REG_CUENTA: índice UNIQUE sobre (subsistema_id, nro_cuenta) — distingue cliente y vendedor con el mismo nro_cuenta.
  • REG_CUENTA_RELACIONES: índice UNIQUE sobre (regcuenta_id, regcuentarelacionada_id, subsistemarelacionada_id) — hace idempotente el upsert de la relación.

delete() en multi-op

MysqlAdapter.delete(entity, pk) solo borra de la primera tabla declarada en operations (la "cabeza" lógica de la entity), usando la primera columna de su upsertKey como filtro. Las demás filas relacionadas no se tocan automáticamente — si hace falta borrado en cascada hay que configurarlo a nivel de base (FK con ON DELETE CASCADE).

MSSQL Destination Adapter

Escribe datos a SQL Server (Tango) mediante queries personalizadas.

Configuración

// src/config/destinations/tango.ts
export const tangoAdapter: MssqlDestAdapterConfig = {
  name: 'tango',
  type: 'mssql',
  config: {
    server: process.env.MSSQL_SERVER_SAVE || 'localhost',
    port: parseInt(process.env.MSSQL_PORT_SAVE || '1433'),
    database: process.env.MSSQL_DATABASE_SAVE || 'tango',
    user: process.env.MSSQL_USER_SAVE || 'sa',
    password: process.env.MSSQL_PASSWORD_SAVE || '',
  },
  entities: {
    pedidosPreparados: {
      transform: 'procesarPedidoPreparado',
      queries: [
        {
          sql: `UPDATE gva21 
                SET aprueba='WMS', estado=2, fecha_apru={{fechaSQL}}
                WHERE nro_pedido={{idERP}}`,
        },
        {
          sql: `INSERT INTO gva77 (...) VALUES (...)`,
          condition: (data) => data.estado === 'PREPARADO',
        },
      ],
      notifyOnSuccess: {
        url: 'http://wms/pedidos/preparados/sincronizar',
        headers: { 'X-API-KEY': '...' },
        bodyBuilder: (data) => ({ pedidos: [data.id] }),
      },
    },
  },
};

Entity Config MSSQL

Propiedad Tipo Descripción
transform string? Transform a aplicar
queries QueryConfig[] Queries a ejecutar
notifyOnSuccess NotifyConfig? Webhook post-commit

Query Config

{
  sql: string;              // Query con placeholders {{campo}}
  condition?: (data) => boolean;  // Ejecutar solo si true
}

Interpolación de Queries

Los placeholders {{campo}} se reemplazan con valores escapados:

// Template
"UPDATE tabla SET col={{valor}} WHERE id={{id}}"

// Data
{ valor: "O'Brien", id: 123 }

// Resultado
"UPDATE tabla SET col='O''Brien' WHERE id=123"

Webhook (notifyOnSuccess)

Después de un commit exitoso, puede notificar a un servicio externo:

sequenceDiagram
    participant Worker
    participant MSSQL
    participant Webhook

    Worker->>MSSQL: BEGIN TRANSACTION
    Worker->>MSSQL: Query 1
    Worker->>MSSQL: Query 2
    Worker->>MSSQL: COMMIT

    alt notifyOnSuccess configurado
        Worker->>Webhook: POST /sincronizar
        Webhook-->>Worker: 200 OK
    end

    Worker-->>Queue: Job completado

Mapeo de Campos

Todos los adapters soportan mapeo de campos vía fields:

fields: {
  // campo_destino: 'CAMPO_ORIGEN'
  codigo: 'CODIGO',
  nombre: 'DESCRIPCION',
  precioVenta: 'PRECIO_LISTA_1',
}

El mapData del BaseAdapter transforma:

// Input
{ CODIGO: 'ART001', DESCRIPCION: 'Producto', PRECIO_LISTA_1: 100 }

// Output
{ codigo: 'ART001', nombre: 'Producto', precioVenta: 100 }

Agregar Nuevo Destination

  1. Crear clase que extienda BaseAdapter
  2. Implementar send() y delete()
  3. Crear archivo de configuración
  4. Agregar al registry por tipo
// Ejemplo: PostgreSQL adapter
export class PostgresAdapter extends BaseAdapter implements Adapter {
  async send(entity: string, pk: string, data: RowData): Promise<AdapterResult> {
    // Implementación
    return this.success();
  }

  async delete(entity: string, pk: string): Promise<AdapterResult> {
    // Implementación
    return this.success();
  }
}