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
Ejemplo:
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
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()