import { pool } from "../db/pool.ts"; import { ctxDe } from "./ctx.ts"; import { buscarConversaciones, mensajesDeConversacion, obtenerConversacion, normalizarTipo, type CrmConversation, type CrmMessage, } from "./conversations.ts"; /** * Fecha del CRM → `Date`. * * La API mezcla formatos: `lastMessageDate` llega como epoch en milisegundos y * `dateAdded` como ISO. Aceptar los dos aquí evita repartir esa comprobación por * todos los sitios que guardan una fecha. */ function fecha(v: string | number | undefined | null): Date | null { if (v == null) return null; const d = new Date(v); return isNaN(d.getTime()) ? null : d; } export async function upsertConversacion( businessId: number, c: CrmConversation ): Promise { const { rows } = await pool.query<{ id: number }>( `INSERT INTO conversations (business_id, crm_conversation_id, crm_contact_id, contact_name, last_message_type, last_message_body, last_message_at, unread_count, client_id, synced_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8, (SELECT id FROM clients WHERE business_id = $1 AND crm_contact_id = $3 AND deleted_at IS NULL LIMIT 1), now()) ON CONFLICT (business_id, crm_conversation_id) DO UPDATE SET crm_contact_id = COALESCE(EXCLUDED.crm_contact_id, conversations.crm_contact_id), -- Ni el nombre ni el canal se degradan. -- -- MEDIDO: GET /conversations/{id} NO devuelve el nombre del contacto ni -- un canal reconocible; eso solo viene del buscador. Sincronizar un hilo -- por su id sobrescribia un nombre bueno con "Sin nombre" y el canal con -- "Desconocido". Un dato pobre no puede pisar a uno que ya se tenia. contact_name = CASE WHEN EXCLUDED.contact_name = 'Sin nombre' THEN COALESCE(conversations.contact_name, EXCLUDED.contact_name) ELSE EXCLUDED.contact_name END, last_message_type = CASE WHEN EXCLUDED.last_message_type = 'Desconocido' THEN COALESCE(conversations.last_message_type, EXCLUDED.last_message_type) ELSE EXCLUDED.last_message_type END, last_message_body = COALESCE(EXCLUDED.last_message_body, conversations.last_message_body), last_message_at = COALESCE(EXCLUDED.last_message_at, conversations.last_message_at), unread_count = EXCLUDED.unread_count, -- El enlace con la clienta solo se RELLENA, nunca se borra: si la -- sincronizacion de contactos todavia no ha corrido, client_id es NULL, -- y pisarlo con NULL mas tarde perderia un enlace ya resuelto. client_id = COALESCE(conversations.client_id, EXCLUDED.client_id), synced_at = now() RETURNING id`, [ businessId, c.id, c.contactId ?? null, c.fullName || c.contactName || "Sin nombre", normalizarTipo(c.lastMessageType), c.lastMessageBody ?? null, fecha(c.lastMessageDate), c.unreadCount ?? 0, ] ); return rows[0].id; } export async function upsertMensaje( businessId: number, conversationId: number, m: CrmMessage ): Promise { await pool.query( `INSERT INTO messages (business_id, conversation_id, crm_message_id, crm_contact_id, direction, channel, channel_raw, body, status, sent_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10) ON CONFLICT (business_id, crm_message_id) DO UPDATE SET -- El CRM es el dueno del historico: aqui se reescribe desde el, nunca se -- edita. Solo cambian cuerpo y estado; el resto es inmutable. body = EXCLUDED.body, status = EXCLUDED.status`, [ businessId, conversationId, m.id, m.contactId ?? null, m.direction === "outbound" ? "outbound" : "inbound", normalizarTipo(m.messageType), m.messageType ?? null, m.body ?? null, m.status ?? null, fecha(m.dateAdded), ] ); } export interface ResumenSyncConv { conversaciones: number; mensajes: number; runId: number; } /** * Espeja UNA conversación con todos sus mensajes. * * Pagina hasta 20 vueltas de 100: son 2 000 mensajes por hilo, muy por encima de * cualquier conversación real, y el tope existe para que un `nextPage` que nunca * deje de ser `true` no cuelgue la petición para siempre. */ export async function sincronizarConversacion( businessId: number, crmConversationId: string ): Promise<{ conversacion: number; mensajes: number }> { const ctx = await ctxDe(businessId); const c = await obtenerConversacion(ctx, crmConversationId); if (!c) throw { status: 404, error: "Esa conversación no existe en Bucéfalo CRM" }; // El endpoint de una conversación suelta no trae el nombre del contacto. Si // la clienta ya está en la plataforma, se usa el suyo: es mejor dato que el // relleno, y evita que la bandeja muestre "Sin nombre" para alguien conocido. let nombre = c.fullName || c.contactName; if (!nombre && c.contactId) { const { rows } = await pool.query<{ name: string }>( `SELECT name FROM clients WHERE business_id = $1 AND crm_contact_id = $2 AND deleted_at IS NULL LIMIT 1`, [businessId, c.contactId] ); nombre = rows[0]?.name; } const convId = await upsertConversacion(businessId, { ...c, id: crmConversationId, fullName: nombre, }); let cursor: string | undefined; let total = 0; for (let i = 0; i < 20; i++) { const { mensajes, lastMessageId, hayMas } = await mensajesDeConversacion( ctx, crmConversationId, { limit: 100, lastMessageId: cursor } ); for (const m of mensajes) { await upsertMensaje(businessId, convId, m); total++; } if (!hayMas || !lastMessageId || !mensajes.length) break; cursor = lastMessageId; } // El canal del hilo se deduce de su ultimo mensaje real. // // GET /conversations/{id} no devuelve un canal reconocible, pero los mensajes // que acabamos de traer si lo traen. Deducirlo de ahi es mejor que dejar // "Desconocido" en la bandeja, y no cuesta ni una peticion mas. // Se excluyen las actividades: son notas que el propio CRM escribe en el hilo, // no un canal por el que hablar con la clienta. await pool.query( `UPDATE conversations c SET last_message_type = COALESCE( (SELECT m.channel FROM messages m WHERE m.conversation_id = c.id AND m.channel <> 'Actividad' ORDER BY m.sent_at DESC NULLS LAST, m.id DESC LIMIT 1), c.last_message_type) WHERE c.id = $1 AND c.last_message_type = 'Desconocido'`, [convId] ); return { conversacion: convId, mensajes: total }; } /** Espeja las conversaciones más recientes con sus últimos mensajes. */ export async function sincronizarConversaciones( businessId: number, opts: { limit?: number; userId?: number | null } = {} ): Promise { const ctx = await ctxDe(businessId); const { rows: run } = await pool.query<{ id: number }>( `INSERT INTO crm_sync_runs (business_id, kind, direction, started_by_user_id) VALUES ($1,'conversations','pull',$2) RETURNING id`, [businessId, opts.userId ?? null] ); const runId = run[0].id; try { const { conversations } = await buscarConversaciones(ctx, { limit: opts.limit ?? 50 }); let mensajes = 0; for (const c of conversations) { const convId = await upsertConversacion(businessId, c); const { mensajes: ms } = await mensajesDeConversacion(ctx, c.id, { limit: 50 }); for (const m of ms) { await upsertMensaje(businessId, convId, m); mensajes++; } } await pool.query( `UPDATE crm_sync_runs SET status='ok', finished_at=now(), fetched=$2, created=$3 WHERE id = $1`, [runId, conversations.length, mensajes] ); return { conversaciones: conversations.length, mensajes, runId }; } catch (e: any) { await pool.query( `UPDATE crm_sync_runs SET status='error', finished_at=now(), error=$2 WHERE id=$1`, [runId, String(e?.error ?? e?.message ?? e).slice(0, 500)] ); throw e; } }