Saltar a contenido

Servicios Core

El middleware tiene cuatro servicios principales que orquestan la sincronización.

Diagrama de Servicios

flowchart TB
    subgraph Services
        POLLER[Poller]
        HASH[HashStore]
        QUEUE[JobQueue]
        METRICS[MetricsStore]
    end

    subgraph External
        REDIS[(Redis)]
    end

    POLLER --> HASH
    POLLER --> QUEUE
    POLLER --> METRICS
    HASH --> REDIS
    QUEUE --> REDIS
    METRICS --> REDIS

Poller

El Poller es el corazón del sistema. Ejecuta el ciclo de polling para cada flow.

Responsabilidades

  • Leer datos de sources
  • Aplicar transforms
  • Detectar cambios (comparar hashes)
  • Encolar jobs
  • Manejar pause/resume/limits

Estado

class Poller {
  private timers: Map<string, NodeJS.Timeout>;  // Timers por flow
  private polling: Set<string>;                  // Flows en ejecución
  private paused: Set<string>;                   // Flows pausados
  private runLimits: Map<string, number>;        // Límites de ejecución
  private runCounts: Map<string, number>;        // Contadores
  private lastSync: Map<string, Date>;           // Última sync (incremental)
}

Ciclo de Polling

flowchart TD
    A[Inicio Poll] --> B{¿Pausado?}
    B -->|Sí| Z[Skip]
    B -->|No| C{¿Límite alcanzado?}
    C -->|Sí| D[Auto-pause]
    C -->|No| E{¿Ya en progreso?}
    E -->|Sí| Z
    E -->|No| F[Leer Source]
    F --> G[Aplicar Transform]
    G --> H[Comparar Hashes]
    H --> I[Encolar Jobs]
    I --> J[Guardar Hashes]
    J --> K[Registrar Métricas]
    K --> L[Incrementar Contador]
    D --> Z

Métodos Principales

// Iniciar todos los flows
startAll(flows: SyncFlowConfig[]): void

// Ejecutar poll de un flow
pollFlow(flow: SyncFlowConfig): Promise<number>

// Ejecutar poll de un solo item
pollOne(flow: SyncFlowConfig, pk: string, force?: boolean): Promise<{
  found: boolean;
  changed: boolean;
  jobs: number;
}>

// Control de pause
pauseFlow(name: string): Promise<void>
resumeFlow(name: string): Promise<void>
pauseAll(flowNames: string[]): Promise<void>
resumeAll(): Promise<void>

// Run limits
setRunLimit(name: string, limit: number): void
clearRunLimit(name: string): void

// Estado
getFlowStatus(name: string): FlowStatus
getAllStatus(): Record<string, FlowStatus>

// Persistencia (Redis)
loadPausedState(): Promise<void>

Persistencia de Estado

El estado de pause se persiste en Redis:

const PAUSE_KEY = 'sync:paused';

// Cargar al iniciar
await poller.loadPausedState();

// Guardar al pausar
await redis.sadd(PAUSE_KEY, flowName);

// Quitar al resumir
await redis.srem(PAUSE_KEY, flowName);

HashStore

Almacena hashes MD5 de cada registro para detectar cambios.

Estructura en Redis

sync:hash:{flow}:{pk} = {hash}

Ejemplo:

sync:hash:articulos:ART001 = "a1b2c3d4e5f6..."
sync:hash:clientes:CLI001 = "f6e5d4c3b2a1..."

Métodos

// Individual
getHash(flow: string, pk: string): Promise<string | null>
setHash(flow: string, pk: string, hash: string): Promise<void>
deleteHash(flow: string, pk: string): Promise<void>

// Batch (más eficiente)
getHashes(flow: string, pks: string[]): Promise<Map<string, string>>
setHashes(flow: string, hashes: Map<string, string>): Promise<void>

// Listado
getAllPKs(flow: string): Promise<string[]>

Uso de Pipeline

Para operaciones batch, usa Redis pipeline:

async setHashes(flow: string, hashes: Map<string, string>): Promise<void> {
  const pipeline = this.redis.pipeline();

  for (const [pk, hash] of hashes) {
    pipeline.set(`sync:hash:${flow}:${pk}`, hash);
  }

  await pipeline.exec();
}

JobQueue

Cola de trabajos basada en BullMQ.

Configuración

const queue = new JobQueue(redis, logger, metrics);

// Worker con concurrencia
queue.startWorker(adapterRegistry, {
  concurrency: 4,
  limiter: { max: 50, duration: 1000 },  // Rate limit
});

Prioridades

Prioridad Valor Descripción
high 1 Se procesa primero
normal 5 Prioridad default
low 10 Se procesa último

Job Structure

interface SyncJob {
  entity: string;    // Nombre del flow
  adapter: string;   // Destination adapter
  action: 'upsert' | 'delete';
  pk: string;
  data: RowData;
}

Retry con Backoff

{
  attempts: 5,
  backoff: {
    type: 'exponential',
    delay: 5000,  // 5s, 10s, 20s, 40s, 80s
  },
}

Degradación de Prioridad

Jobs fallidos bajan de prioridad:

flowchart LR
    A[Intento 1<br/>High] -->|Falla| B[Intento 2<br/>Normal]
    B -->|Falla| C[Intento 3<br/>Low]
    C -->|Falla| D[Intento 4<br/>Low]
    D -->|Falla| E[Intento 5<br/>Low]
    E -->|Falla| F[Fallo Permanente]

Worker Events

worker.on('completed', (job) => {
  // Job completado
});

worker.on('failed', (job, error) => {
  // Job falló - se reintentará si tiene attempts
});

worker.on('error', (error) => {
  // Error del worker
});

MetricsStore

Almacena métricas de cada flow en Redis.

Datos Almacenados

interface FlowMetrics {
  lastRun: string;           // ISO timestamp
  avgDuration: number;       // ms promedio
  p95Duration: number;       // Percentil 95
  totalRuns: number;
  totalErrors: number;
  lastRunDetails: {
    timestamp: string;
    durationMs: number;
    rowsFound: number;
    changes: number;
    jobsCreated: number;
    timings: {
      query: number;
      transform: number;
      compareHash: number;
      detectDeletes: number;
      enqueue: number;
      saveHashes: number;
    };
    error?: string;
  };
}

Métodos

// Registrar ejecución
record(flow: string, details: RunDetails): Promise<void>

// Obtener métricas
get(flow: string): Promise<FlowMetrics | null>
getAll(): Promise<Record<string, FlowMetrics>>

Cálculo de P95

Mantiene historial de últimas 100 ejecuciones para calcular percentil 95:

const sorted = durations.sort((a, b) => a - b);
const p95Index = Math.floor(sorted.length * 0.95);
const p95 = sorted[p95Index];

Interacción entre Servicios

sequenceDiagram
    participant API
    participant Poller
    participant HashStore
    participant Queue
    participant Worker
    participant Metrics
    participant Adapter

    API->>Poller: POST /sync/articulos
    Poller->>Poller: pollFlow(articulos)
    Poller->>HashStore: getHashes(pks)
    HashStore-->>Poller: existingHashes

    loop Por cada cambio
        Poller->>Queue: addJob(upsert)
    end

    Poller->>HashStore: setHashes(newHashes)
    Poller->>Metrics: record(flow, details)

    Worker->>Queue: getJob()
    Queue-->>Worker: job
    Worker->>Adapter: send(data)
    Adapter-->>Worker: result
    Worker->>Metrics: updateStats()