diff --git a/apps/api/.env.example b/apps/api/.env.example index 6ffd2fc..28e5b26 100644 --- a/apps/api/.env.example +++ b/apps/api/.env.example @@ -101,3 +101,6 @@ SAT_PROXY_LIST= SAT_PROXY_STRATEGY=round-robin # Si es true, cuando todos los proxies fallan se intenta con la IP directa del servidor. SAT_PROXY_FALLBACK_DIRECT=true +# Máximo de contribuyentes sincronizados en paralelo por el scheduler (default: 10). +# Aprovecha el pool de proxies; cada RFC suele usar una IP distinta en round-robin. +SAT_CONCURRENT_CONTRIBUYENTES=10 diff --git a/apps/api/src/jobs/sat-sync.job.ts b/apps/api/src/jobs/sat-sync.job.ts index 1449ea7..9c9e66f 100644 --- a/apps/api/src/jobs/sat-sync.job.ts +++ b/apps/api/src/jobs/sat-sync.job.ts @@ -15,7 +15,8 @@ const SYNC_CRON_SCHEDULE = '0 6-10 * * *'; // 6:00–10:00 AM CDMX — ~20% de t const RECOVERY_CRON_SCHEDULE = '0 10 * * *'; // 10:00 AM todos los días const RETRY_9AM_CRON_SCHEDULE = '0 9 * * *'; // 9:00 AM todos los días const RETRY_4PM_CRON_SCHEDULE = '0 16 * * *'; // 4:00 PM todos los días -const CONCURRENT_SYNCS = 3; // Máximo de sincronizaciones simultáneas +const CONCURRENT_SYNCS = 3; // Máximo de sincronizaciones simultáneas (legacy, se mantiene por compatibilidad) +const CONCURRENT_CONTRIBUYENTES = Number(process.env.SAT_CONCURRENT_CONTRIBUYENTES || '10'); // Máximo de contribuyentes en paralelo const OPINION_CRON_SCHEDULE = '0 4 * * 0'; // Sundays 4:00 AM const CSF_CRON_SCHEDULE = '0 4 1 * *'; // Día 1 de cada mes 04:00 AM (CSF mensual) const INCREMENTAL_CRON_SCHEDULE = '0 11,15,19 * * *'; // 11:00, 15:00 y 19:00; fuera de ese rango el daily (6-10 AM) cubre @@ -133,7 +134,100 @@ async function getContribuyentesParaSync( } /** - * Ejecuta sincronización para un tenant y sus contribuyentes + * Unidad mínima de sincronización: un tenant (legacy) o un contribuyente. + */ +interface SyncUnit { + tenantId: string; + contribuyenteId?: string; + syncType: 'initial' | 'daily' | 'incremental'; +} + +/** + * Recolecta todas las unidades de sync para un conjunto de tenants. + * - Modo incremental: solo incluye contribuyentes/tenants con initial completado. + * - Modo daily/initial: determina initial vs daily por contribuyente. + */ +async function getSyncUnits( + tenantIds: string[], + options: { incremental?: boolean; logPrefix?: string } = {} +): Promise { + const { incremental = false, logPrefix = '[SAT Cron]' } = options; + const units: SyncUnit[] = []; + + for (const tenantId of tenantIds) { + try { + const tenant = await prisma.tenant.findUnique({ + where: { id: tenantId }, + select: { databaseName: true }, + }); + + let contribuyenteIds: string[] = []; + if (tenant?.databaseName) { + const { ids, total } = await getContribuyentesParaSync(tenantId, tenant.databaseName, logPrefix); + if (total > 0 && ids.length === 0) { + console.log(`${logPrefix} Tenant ${tenantId}: ningún contribuyente con FIEL vigente, se omite`); + continue; + } + contribuyenteIds = ids; + } + + // Tenant legacy sin contribuyentes (Horux 360) + if (contribuyenteIds.length === 0) { + if (incremental) { + const hasInitial = await prisma.satSyncJob.findFirst({ + where: { tenantId, contribuyenteId: null, type: 'initial', status: 'completed' }, + }); + if (!hasInitial) continue; + units.push({ tenantId, syncType: 'incremental' }); + } else { + const needsInitial = await needsInitialSync(tenantId); + units.push({ tenantId, syncType: needsInitial ? 'initial' : 'daily' }); + } + continue; + } + + // Contribuyentes del tenant + for (const contribuyenteId of contribuyenteIds) { + if (incremental) { + const hasInitial = await prisma.satSyncJob.findFirst({ + where: { tenantId, contribuyenteId, type: 'initial', status: 'completed' }, + }); + if (!hasInitial) continue; + units.push({ tenantId, contribuyenteId, syncType: 'incremental' }); + } else { + const needsInitial = await needsInitialSync(tenantId, contribuyenteId); + units.push({ tenantId, contribuyenteId, syncType: needsInitial ? 'initial' : 'daily' }); + } + } + } catch (error: any) { + console.error(`${logPrefix} Error recolectando unidades para tenant ${tenantId}:`, error.message); + } + } + + return units; +} + +/** + * Ejecuta sync para una unidad (tenant o contribuyente), respetando locks. + */ +async function syncUnit(unit: SyncUnit, logPrefix: string): Promise { + try { + const status = await getSyncStatus(unit.tenantId, unit.contribuyenteId); + if (status.hasActiveSync) { + console.log(`${logPrefix} ${unit.tenantId}${unit.contribuyenteId ? ` contribuyente ${unit.contribuyenteId}` : ''} ya tiene sync activo, omitiendo`); + return; + } + + console.log(`${logPrefix} Iniciando sync ${unit.syncType} para ${unit.tenantId}${unit.contribuyenteId ? ` contribuyente ${unit.contribuyenteId}` : ''}`); + const jobId = await startSync(unit.tenantId, unit.syncType, undefined, undefined, unit.contribuyenteId); + console.log(`${logPrefix} Job ${jobId} iniciado`); + } catch (error: any) { + console.error(`${logPrefix} Error sincronizando ${unit.tenantId}${unit.contribuyenteId ? ` contribuyente ${unit.contribuyenteId}` : ''}:`, error.message); + } +} + +/** + * Ejecuta sincronización para un tenant y sus contribuyentes (modo secuencial legacy) */ async function syncTenant(tenantId: string): Promise { try { @@ -247,13 +341,22 @@ async function runSyncJob(): Promise { return; } - // Procesar en lotes para no saturar - for (let i = 0; i < groupTenants.length; i += CONCURRENT_SYNCS) { - const batch = groupTenants.slice(i, i + CONCURRENT_SYNCS); - await Promise.all(batch.map(syncTenant)); + // Recolectar unidades de sync (contribuyentes o tenants legacy) + const units = await getSyncUnits(groupTenants, { logPrefix: '[SAT Cron]' }); + console.log(`[SAT Cron] Ventana ${hour}:00 CDMX — ${units.length} unidades de sync listas (max ${CONCURRENT_CONTRIBUYENTES} paralelas)`); + + if (units.length === 0) { + console.log('[SAT Cron] No hay unidades de sync en este grupo'); + return; + } + + // Procesar en lotes de contribuyentes para aprovechar los proxies + for (let i = 0; i < units.length; i += CONCURRENT_CONTRIBUYENTES) { + const batch = units.slice(i, i + CONCURRENT_CONTRIBUYENTES); + await Promise.all(batch.map(unit => syncUnit(unit, '[SAT Cron]'))); // Pequeña pausa entre lotes - if (i + CONCURRENT_SYNCS < groupTenants.length) { + if (i + CONCURRENT_CONTRIBUYENTES < units.length) { await new Promise(resolve => setTimeout(resolve, 5000)); } } @@ -385,11 +488,19 @@ async function runIncrementalSyncJob(): Promise { if (tenantIds.length === 0) return; - for (let i = 0; i < tenantIds.length; i += CONCURRENT_SYNCS) { - const batch = tenantIds.slice(i, i + CONCURRENT_SYNCS); - await Promise.all(batch.map(incrementalSyncTenant)); + const units = await getSyncUnits(tenantIds, { incremental: true, logPrefix: '[SAT Cron Inc]' }); + console.log(`[SAT Cron Inc] ${units.length} unidades de sync listas (max ${CONCURRENT_CONTRIBUYENTES} paralelas)`); - if (i + CONCURRENT_SYNCS < tenantIds.length) { + if (units.length === 0) { + console.log('[SAT Cron Inc] No hay unidades de sync'); + return; + } + + for (let i = 0; i < units.length; i += CONCURRENT_CONTRIBUYENTES) { + const batch = units.slice(i, i + CONCURRENT_CONTRIBUYENTES); + await Promise.all(batch.map(unit => syncUnit(unit, '[SAT Cron Inc]'))); + + if (i + CONCURRENT_CONTRIBUYENTES < units.length) { await new Promise(resolve => setTimeout(resolve, 5000)); } } diff --git a/docs/SAT-SYNC-IMPLEMENTATION.md b/docs/SAT-SYNC-IMPLEMENTATION.md index 5aa9056..9f82b8d 100644 --- a/docs/SAT-SYNC-IMPLEMENTATION.md +++ b/docs/SAT-SYNC-IMPLEMENTATION.md @@ -17,6 +17,7 @@ Los datos se almacenan en la base de datos del tenant correspondiente. | `apps/api/src/services/sat/sat.service.ts` | Lógica principal de sincronización, políticas de reintento, polling | | `apps/api/src/services/sat/sat-client.service.ts` | Cliente del SAT, `query`, `verify`, `download` | | `apps/api/src/services/sat/sat-parser.service.ts` | Parseo de XMLs y metadata | +| `apps/api/src/services/sat/proxy.service.ts` | Pool rotativo de proxies SAT | | `apps/api/src/services/sat/sat-crypto.service.ts` | Encriptación AES-256-GCM de credenciales FIEL | | `apps/api/src/services/fiel.service.ts` | FIEL a nivel tenant (legacy) | | `apps/api/src/services/contribuyente-fiel.service.ts` | FIEL por contribuyente (modelo despacho) | @@ -87,7 +88,7 @@ Definidos en `apps/api/src/jobs/sat-sync.job.ts`: | Job | Expresión | Horario CDMX | Propósito | |-----|-----------|--------------|-----------| -| SAT Cron | `0 6-10 * * *` | 6:00–10:00 AM | Daily sync, ~20% de tenants por hora | +| SAT Cron | `0 6-10 * * *` | 6:00–10:00 AM | Daily sync, ~20% de tenants por hora, hasta `SAT_CONCURRENT_CONTRIBUYENTES` paralelos | | Recovery Cron | `0 10 * * *` | 10:00 AM | Recuperar jobs `running` atorados | | Daily Retry | `0 9 * * *` y `0 16 * * *` | 9:00 AM y 4:00 PM | Reintentar daily fallidos | | Incremental Enterprise | `0 11,15,19 * * *` | 11 AM, 3 PM, 7 PM | Sync incremental | @@ -116,7 +117,46 @@ Definidos en `apps/api/src/jobs/sat-sync.job.ts`: 1. Ventana de 8 horas: `ahora - 10h` a `ahora - 2h`. 2. Descarga XMLs + metadata de emitidos y recibidos. -## 7. Polling y límites +## 7. Proxies SAT (rotación por IP) + +Para mitigar el bloqueo `404 Error no controlado` causado por cuota de solicitudes desde una sola IP pública, el sistema soporta un pool de proxies HTTP/HTTPS rotativos. + +### Configuración + +Variables en `apps/api/.env`: + +```bash +# Lista de proxies separados por coma. Soporta autenticación básica. +SAT_PROXY_LIST=http://user:pass@host1:port,http://user:pass@host2:port + +# Estrategia de rotación: round-robin | random (default: round-robin) +SAT_PROXY_STRATEGY=round-robin + +# Si true, cuando todos los proxies fallan se intenta con la IP directa del servidor. +SAT_PROXY_FALLBACK_DIRECT=true + +# Máximo de contribuyentes sincronizados en paralelo por el scheduler (default: 10). +SAT_CONCURRENT_CONTRIBUYENTES=10 +``` + +### Componentes + +- **`ProxyManager`** (`apps/api/src/services/sat/proxy.service.ts`): parsea `SAT_PROXY_LIST`, rota proxies y crea `HttpsProxyAgent`. +- **`sat-client.service.ts`**: usa el agente del proxy en cada petición SOAP al SAT. + +### Comportamiento + +- Cada llamada a `getNextProxy()` devuelve el siguiente proxy del pool (round-robin). +- Cuando el scheduler lanza 10 contribuyentes en paralelo, cada uno tiende a usar una IP distinta. +- Si un proxy devuelve error de conexión, el cliente SAT puede caer a la IP directa según `SAT_PROXY_FALLBACK_DIRECT`. + +### Recomendaciones operativas + +- Tamaño mínimo del pool: **1 proxy por cada 5 RFCs** que se sincronicen en paralelo. +- Con `SAT_CONCURRENT_CONTRIBUYENTES=10`, un pool de 10 proxies da una IP por RFC en el peor caso. +- Monitorear logs por `[ProxyManager] Usando proxy: ...` y por 404 persistentes en una misma IP. + +## 8. Polling y límites Después de crear una solicitud (`query`) al SAT, el sistema verifica el estado periódicamente (`verify`). @@ -133,7 +173,7 @@ Esto da un máximo de **~45 minutos por solicitud** (9 × 5 min). Cada solicitud al SAT tiene su propio polling; los intentos no se comparten entre solicitudes. -## 8. Políticas de reintentos +## 9. Políticas de reintentos ```ts const RETRY_POLICIES = { @@ -148,7 +188,7 @@ const MAX_DAILY_RETRY_ATTEMPTS = 5; // original + 2 automáticos + 2 crons fijos Los reintentos automáticos se programan desde `createdAt` del job. Los reintentos por cron fijo (9 AM / 4 PM) se manejan en `continuePendingDailyRequests`. -## 9. Manejo de errores +## 10. Manejo de errores ### Errores no fatales (se registran, no abortan en daily) @@ -166,7 +206,7 @@ Los reintentos automáticos se programan desde `createdAt` del job. Los reintent - FIEL inválida o vencida. - Errores que no son transitorios y no están en la lista de no fatales. -## 10. Errores comunes del SAT +## 11. Errores comunes del SAT | Código/Mensaje | Significado | Acción | |----------------|-------------|--------| @@ -178,7 +218,7 @@ Los reintentos automáticos se programan desde `createdAt` del job. Los reintent | "Fecha final invalida" | Fecha futura o mal formada | Usar `getYesterdayEnd()` | | "El certificado no es válido" | FIEL rechazada por el SAT | Revisar vigencia/contraseña de FIEL | -## 11. Monitoreo y comandos útiles +## 12. Monitoreo y comandos útiles ### Estado del API @@ -223,7 +263,13 @@ psql "$DATABASE_URL" -c " " ``` -## 12. Changelog reciente +## 13. Changelog reciente + +### 2026-08-03 + +- **Proxies SAT rotativos**: integración de `ProxyManager` y pool de proxies HTTP/HTTPS para evitar bloqueo por IP del SAT. +- **Concurrencia por contribuyente**: el scheduler daily e incremental procesa hasta `SAT_CONCURRENT_CONTRIBUYENTES=10` RFCs en paralelo, sin importar a cuántos tenants pertenezcan. +- Variables de entorno: `SAT_PROXY_LIST`, `SAT_PROXY_STRATEGY`, `SAT_PROXY_FALLBACK_DIRECT`, `SAT_CONCURRENT_CONTRIBUYENTES`. ### 2026-07-31 @@ -233,14 +279,15 @@ psql "$DATABASE_URL" -c " - Daily retry fijo a las 9:00 AM y 4:00 PM CDMX. - Incremental Enterprise a las 11:00 AM, 3:00 PM y 7:00 PM CDMX. -## 13. Problemas conocidos +## 14. Problemas conocidos -1. **Bloqueo `404 Error no controlado` del SAT**: Aparece cuando se hacen muchas consultas desde la misma IP. Mitigación temporal: reducir frecuencia de polling y metadata solo domingos. +1. **Bloqueo `404 Error no controlado` del SAT**: Aparece cuando se hacen muchas consultas desde la misma IP. Mitigación: proxies rotativos (`SAT_PROXY_LIST`) + concurrencia controlada por contribuyente + metadata histórica solo domingos. 2. **FIEL inválida**: Algunos tenants/contribuyentes tienen FIEL rechazada por el SAT. Requiere revisar/renovar FIEL. 3. **Jobs `initial` atorados en `running`**: Pueden quedar si el proceso se reinicia; el recovery cron y el watchdog los limpian. -## 14. Próximos pasos +## 15. Próximos pasos -- [ ] Implementar proxies rotativos para evitar bloqueo por IP del SAT. -- [ ] Monitorear tasa de éxito tras reducir polling y metadata solo domingos. +- [x] Implementar proxies rotativos para evitar bloqueo por IP del SAT. +- [ ] Monitorear tasa de éxito tras proxies + concurrencia por contribuyente. - [ ] Revisar/renovar FIELs inválidas reportadas por el monitor. +- [ ] Evaluar ampliar ventana horaria del daily (6–10 AM) si el volumen de RFCs supera el throughput con 10 paralelos.