Handler completo para CULTIVA Visuals: verifica firma, previene replay attacks
y enruta eventos del ciclo de vida de predicciones hacia Supabase + Ably.
// ============================================================ // webhook-replicate.js โ CULTIVA Visuals // Handler verificado para webhooks de Replicate // Integra Supabase (persistencia) + Ably (notificaciones RT) // ============================================================ const express = require('express'); const crypto = require('crypto'); const { createClient } = require('@supabase/supabase-js'); const Ably = require('ably'); const app = express(); // โโ Clientes de servicios externos โโโโโโโโโโโโโโโโโโโโโโโโโโ const supabase = createClient( process.env.SUPABASE_URL, process.env.SUPABASE_SERVICE_ROLE_KEY ); const ably = new Ably.Realtime({ key: process.env.ABLY_API_KEY }); // โโ Verificacion de firma HMAC-SHA256 โโโโโโโโโโโโโโโโโโโโโโโ function verificarFirmaReplicate(body, headers, secret) { const { 'webhook-id': id, 'webhook-timestamp': ts, 'webhook-signature': sig } = headers; if (!id || !ts || !sig) throw new Error('Cabeceras webhook obligatorias ausentes'); // Extraer clave: eliminar prefijo 'whsec_' y decodificar base64 const clave = Buffer.from(secret.split('_')[1], 'base64'); const bodyStr = Buffer.isBuffer(body) ? body.toString() : body; const contenido = `${id}.${ts}.${bodyStr}`; const firmaEsperada = crypto .createHmac('sha256', clave) .update(contenido) .digest('base64'); // Replicate puede enviar multiples firmas separadas por espacio const firmas = sig.split(' ').map(s => { const partes = s.split(','); return partes.length > 1 ? partes[1] : s; }); const esValida = firmas.some(f => { try { return crypto.timingSafeEqual(Buffer.from(f), Buffer.from(firmaEsperada)); } catch { return false; } }); // Anti-replay: rechazar si el timestamp supera 5 minutos const ahora = Math.floor(Date.now() / 1000); if (ahora - parseInt(ts, 10) > 300) throw new Error('Timestamp caducado (replay attack)'); return esValida; } // โโ Helpers de negocio โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ async function actualizarPrediccion(prediction) { const { error } = await supabase .from('predictions') .update({ status: prediction.status, output_url: Array.isArray(prediction.output) ? prediction.output[0] : null, error_msg: prediction.error ?? null, duration_s: prediction.metrics?.predict_time ?? null, updated_at: new Date().toISOString(), }) .eq('replicate_id', prediction.id); if (error) throw error; } async function notificarCliente(userId, payload) { const canal = ably.channels.get(`user:${userId}`); await canal.publish('prediction-update', payload); } // โโ Endpoint webhook โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ // CRITICO: express.raw() โ Replicate necesita el body sin parsear app.post('/webhooks/replicate', express.raw({ type: 'application/json' }), async (req, res) => { const t0 = Date.now(); try { // 1. Verificar firma const secret = process.env.REPLICATE_WEBHOOK_SECRET; if (!secret) throw new Error('REPLICATE_WEBHOOK_SECRET no configurado'); const ok = verificarFirmaReplicate(req.body, req.headers, secret); if (!ok) return res.status(400).json({ error: 'Firma invalida' }); // 2. Parsear payload const prediction = JSON.parse(req.body.toString()); const { id, status, version } = prediction; console.log(JSON.stringify({ nivel: 'info', msg: 'webhook recibido', id, status, version })); // 3. Enrutar por estado switch (status) { case 'starting': await actualizarPrediccion(prediction); break; case 'processing': await actualizarPrediccion(prediction); // Emitir logs en tiempo real si los hay if (prediction.logs) { const { data: job } = await supabase .from('predictions').select('user_id').eq('replicate_id', id).single(); if (job) await notificarCliente(job.user_id, { id, status, logs: prediction.logs }); } break; case 'succeeded': await actualizarPrediccion(prediction); { const { data: job } = await supabase .from('predictions').select('user_id').eq('replicate_id', id).single(); if (job) await notificarCliente(job.user_id, { id, status, outputUrl: Array.isArray(prediction.output) ? prediction.output[0] : null, duration: prediction.metrics?.predict_time, }); } break; case 'failed': await actualizarPrediccion(prediction); console.error(JSON.stringify({ nivel: 'error', msg: 'prediccion fallida', id, error: prediction.error })); break; case 'canceled': await actualizarPrediccion(prediction); break; default: console.warn(JSON.stringify({ nivel: 'warn', msg: 'estado desconocido', id, status })); } // 4. Responder rapido (<500ms) para evitar reintentos res.status(200).json({ received: true, status, ms: Date.now() - t0 }); } catch (err) { console.error(JSON.stringify({ nivel: 'error', msg: err.message })); const conocidos = ['Cabeceras webhook obligatorias ausentes', 'Timestamp caducado (replay attack)']; res.status(400).json({ error: conocidos.includes(err.message) ? err.message : 'Webhook invalido' }); } } ); app.use(express.json()); // otras rutas โ DEBE ir despues del endpoint webhook module.exports = { app };
{ "id": "cv_pred_a3f82b1c", "version": "stability-ai/sdxl:39ed52f...", "status": "succeeded", "output": ["https://pbxt.replicate.delivery/cv_out_a3f8.png"], "metrics": { "predict_time": 12.4 }, "created_at": "2026-06-16T09:14:33Z" }
{ "id": "cv_pred_d9e12c44", "version": "black-forest-labs/flux-1.1-pro", "status": "failed", "error": "CUDA out of memory. Tried to allocate 2.5 GiB", "output": null, "created_at": "2026-06-16T09:21:07Z" }
{ "id": "cv_pred_f1a045bb", "status": "processing", "logs": "Loading model... done\nStep 10/50 [ETA 8s]", "output": null, "created_at": "2026-06-16T10:02:14Z" }
{ "id": "cv_pred_8c3e99d2", "status": "starting", "input": { "prompt": "product photo, white background, DSLR", "num_outputs": 1 }, "output": null, "created_at": "2026-06-16T11:30:00Z" }