feat(sat): Fase 6 — concurrencia por contribuyente con proxies rotativos
- Procesa hasta SAT_CONCURRENT_CONTRIBUYENTES=10 RFCs en paralelo - Scheduler daily e incremental trabajan por unidad de sync (contribuyente/tenant legacy) - Agrega SAT_CONCURRENT_CONTRIBUYENTES a .env.example - Actualiza docs/SAT-SYNC-IMPLEMENTATION.md con sección de proxies
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<SyncUnit[]> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
try {
|
||||
@@ -247,13 +341,22 @@ async function runSyncJob(): Promise<void> {
|
||||
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<void> {
|
||||
|
||||
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));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user