import crypto from "node:crypto"; import type { PoolClient } from "pg"; import { pool } from "../db/pool.ts"; import { CrmTransportError } from "./client.ts"; import { proyectarCita } from "./syncAppointments.ts"; export type EntidadOutbox = "appointment" | "client" | "message"; /** * Clave de deduplicación **propia y estable**. Nunca se deriva del contenido: * dos ediciones que dejan el mismo valor son dos intenciones distintas y las * dos tienen que salir. */ export function claveDedup( businessId: number, entidad: EntidadOutbox, entidadId: number, operacion: string, secuencia: number | string ): string { return crypto .createHash("sha256") .update([businessId, entidad, entidadId, operacion, secuencia].join("|")) .digest("hex"); } /** * Encola un cambio para el CRM **dentro de la transacción que lo produjo**. * * Recibe el `PoolClient` a propósito: el cambio local y su fila de bandeja se * escriben juntos o no se escriben. Sin eso aparece la escritura perdida — el * usuario ve «guardado», el proceso muere antes de encolar, y nadie lo reclama * nunca. */ export async function encolar( tx: PoolClient, args: { businessId: number; entidad: EntidadOutbox; entidadId: number; operacion: string; payload: unknown; secuencia?: number | string; } ): Promise { const secuencia = args.secuencia ?? Date.now(); const dedup = claveDedup( args.businessId, args.entidad, args.entidadId, args.operacion, secuencia ); await tx.query( `INSERT INTO crm_outbox (business_id, entity, entity_id, operation, payload, dedup_key) VALUES ($1,$2,$3,$4,$5::jsonb,$6) ON CONFLICT (dedup_key) DO NOTHING`, [ args.businessId, args.entidad, args.entidadId, args.operacion, JSON.stringify(args.payload ?? {}), dedup, ] ); } export interface ResumenDespacho { tomadas: number; confirmadas: number; fallidas: number; indeterminadas: number; } /** * Despacha la bandeja de salida de un negocio. * * FIFO estricto y **una sola escritura en vuelo por registro**: el CRM * estrangula por token y dos escrituras concurrentes sobre la misma cita * corren contra una base que ya cambió. */ export async function despachar( businessId: number, limite = 25 ): Promise { const resumen: ResumenDespacho = { tomadas: 0, confirmadas: 0, fallidas: 0, indeterminadas: 0, }; const { rows } = await pool.query( `SELECT id, entity, entity_id, operation, attempts FROM crm_outbox WHERE business_id = $1 AND status IN ('pendiente','indeterminado') ORDER BY id LIMIT $2`, [businessId, limite] ); resumen.tomadas = rows.length; for (const fila of rows) { await pool.query( `UPDATE crm_outbox SET status = 'enviando', attempts = attempts + 1 WHERE id = $1`, [fila.id] ); try { let crmId: string | null = null; if (fila.entity === "appointment") { const r = await proyectarCita(businessId, fila.entity_id); crmId = r.crmOpportunityId; } else { // Todavía no hay más entidades salientes; se descarta explícitamente // en vez de dejarla girando en la cola para siempre. await pool.query( `UPDATE crm_outbox SET status = 'fallido', last_error = 'entidad no soportada todavía' WHERE id = $1`, [fila.id] ); resumen.fallidas++; continue; } await pool.query( `UPDATE crm_outbox SET status = 'confirmado', crm_id = $2, evidence = 'relectura', sent_at = now(), last_error = NULL WHERE id = $1`, [fila.id, crmId] ); resumen.confirmadas++; } catch (e: any) { // Un fallo de transporte NO se reintenta: el servidor no habló, así que // no se sabe si la escritura entró, y reenviar es fabricar el duplicado. // Queda en `indeterminado` para resolverlo LEYENDO. const indeterminado = e instanceof CrmTransportError || e?.indeterminate === true; await pool.query( `UPDATE crm_outbox SET status = $2, last_error = $3 WHERE id = $1`, [ fila.id, indeterminado ? "indeterminado" : fila.attempts >= 4 ? "fallido" : "pendiente", String(e?.message ?? e).slice(0, 500), ] ); if (indeterminado) resumen.indeterminadas++; else resumen.fallidas++; } } return resumen; } export interface EstadoOutbox { pendiente: number; enviando: number; confirmado: number; fallido: number; indeterminado: number; } export async function estadoOutbox(businessId: number): Promise { const { rows } = await pool.query( `SELECT status, count(*)::int AS c FROM crm_outbox WHERE business_id = $1 GROUP BY status`, [businessId] ); const base: EstadoOutbox = { pendiente: 0, enviando: 0, confirmado: 0, fallido: 0, indeterminado: 0, }; for (const r of rows) (base as any)[r.status] = r.c; return base; }