pushHandler(new RotatingFileHandler(__DIR__ . '/../logs/webhook.log', 7, Logger::INFO)); // Inicializar cola $queue = new RedisQueue(null, $logger); class OptimizedWebhook { private $db; private $logger; private $queue; private $redis; public function __construct($db, $logger, $queue) { $this->db = $db; $this->logger = $logger; $this->queue = $queue; // Inicializar Redis directamente para verificaciones rápidas $this->redis = new Predis\Client([ 'scheme' => getenv('REDIS_SCHEME') ?: 'tcp', 'host' => getenv('REDIS_HOST') ?: '127.0.0.1', 'port' => getenv('REDIS_PORT') ?: 6379, ]); } public function handleRequest() { $method = $_SERVER['REQUEST_METHOD']; if ($method === 'GET') { $this->verifyWebhook(); } elseif ($method === 'POST') { $this->processIncomingMessage(); } else { http_response_code(405); echo json_encode(['error' => 'Método no permitido']); } } private function verifyWebhook() { $verifyToken = $_GET['hub_verify_token'] ?? ''; $challenge = $_GET['hub_challenge'] ?? ''; $mode = $_GET['hub_mode'] ?? ''; if ($mode === 'subscribe' && $verifyToken === WEBHOOK_VERIFY_TOKEN) { $this->logger->info("Webhook verified successfully"); echo $challenge; exit; } $this->logger->warning("Invalid verification attempt", [ 'mode' => $mode, 'token_match' => $verifyToken === WEBHOOK_VERIFY_TOKEN ]); http_response_code(403); echo json_encode(['error' => 'Token de verificación inválido']); } private function processIncomingMessage() { $startTime = microtime(true); // Leer payload $input = file_get_contents('php://input'); $data = json_decode($input, true); // Responder inmediatamente con 200 OK (CRÍTICO: < 5 segundos) http_response_code(200); echo json_encode(['status' => 'received']); // Forzar envío de respuesta al cliente if (function_exists('fastcgi_finish_request')) { fastcgi_finish_request(); } else { // Fallback para otros SAPIs ob_end_flush(); flush(); } // Ahora procesamos el webhook sin presión de tiempo $processingStart = microtime(true); try { // Log webhook (async, no bloquea) $this->logWebhookAsync($input); if (!$data || !isset($data['entry'])) { $this->logger->warning("Invalid webhook payload received"); return; } $this->logger->info("Webhook received", [ 'entries' => count($data['entry']), 'response_time_ms' => round((microtime(true) - $startTime) * 1000, 2) ]); // Procesar cada entrada foreach ($data['entry'] as $entry) { if (!isset($entry['changes'])) continue; foreach ($entry['changes'] as $change) { if (isset($change['value']['messages'])) { $this->queueMessages($change['value']['messages']); } if (isset($change['value']['statuses'])) { $this->updateMessageStatuses($change['value']['statuses']); } } } $totalTime = round((microtime(true) - $startTime) * 1000, 2); $processingTime = round((microtime(true) - $processingStart) * 1000, 2); $this->logger->info("Webhook processed", [ 'total_time_ms' => $totalTime, 'processing_time_ms' => $processingTime, 'response_sent_in_ms' => round(($processingStart - $startTime) * 1000, 2) ]); } catch (Exception $e) { $this->logger->error("Error processing webhook", [ 'error' => $e->getMessage(), 'trace' => $e->getTraceAsString() ]); } } /** * Encolar mensajes para procesamiento asíncrono */ private function queueMessages($messages) { foreach ($messages as $message) { try { $messageId = $message['id'] ?? null; $phoneNumber = $message['from'] ?? null; if (!$messageId || !$phoneNumber) { continue; } // Verificación rápida de duplicados en Redis (más rápido que DB) $duplicateKey = "whatsapp:processed:" . $messageId; if ($this->redis->exists($duplicateKey)) { $this->logger->debug("Duplicate message skipped", ['message_id' => $messageId]); continue; } // Marcar como procesado en Redis (expira en 24h) $this->redis->setex($duplicateKey, 86400, time()); // Obtener o crear usuario (rápido, solo DB lookup) $user = $this->getUserByPhone($phoneNumber); if (!$user) { $userId = $this->createUser($phoneNumber); $user = $this->getUserById($userId); } // Extraer tipo y contenido del mensaje $messageData = $this->extractMessageData($message); // Guardar mensaje en BD (no esperar resultado del bot) $this->saveIncomingMessage( $user['id'], $messageId, $messageData['text'], $messageData['type'], $messageData['media_url'], $messageData['timestamp'] ); // Encolar para procesamiento por worker $this->queue->push('messages', [ 'user' => $user, 'messageText' => $messageData['text'], 'messageType' => $messageData['type'], 'messageId' => $messageId, 'mediaUrl' => $messageData['media_url'] ], 0); // Prioridad alta $this->logger->info("Message queued", [ 'message_id' => $messageId, 'user_id' => $user['id'], 'type' => $messageData['type'] ]); } catch (Exception $e) { $this->logger->error("Failed to queue message", [ 'message_id' => $messageId ?? 'unknown', 'error' => $e->getMessage() ]); } } } /** * Extraer datos del mensaje según tipo */ private function extractMessageData($message) { $type = $message['type'] ?? 'text'; $text = ''; $mediaUrl = null; $timestamp = $message['timestamp'] ?? null; switch ($type) { case 'text': $text = $message['text']['body'] ?? ''; break; case 'interactive': // Normalizar respuestas de botones/listas if (isset($message['interactive']['list_reply']['title'])) { $text = $message['interactive']['list_reply']['title']; } elseif (isset($message['interactive']['button_reply']['title'])) { $text = $message['interactive']['button_reply']['title']; } // Extraer número si viene como "1. Opción" if (preg_match('/^\s*(\d+)\b/', $text, $m)) { $text = $m[1]; } break; case 'image': $text = $message['image']['caption'] ?? ''; $mediaUrl = $message['image']['url'] ?? $message['image']['id'] ?? null; // Encolar descarga de media if ($mediaUrl) { $this->queueMediaDownload($mediaUrl, $message['id'], 'image'); } break; case 'video': $text = $message['video']['caption'] ?? ''; $mediaUrl = $message['video']['url'] ?? $message['video']['id'] ?? null; if ($mediaUrl) { $this->queueMediaDownload($mediaUrl, $message['id'], 'video'); } break; case 'audio': $mediaUrl = $message['audio']['url'] ?? $message['audio']['id'] ?? null; if ($mediaUrl) { $this->queueMediaDownload($mediaUrl, $message['id'], 'audio'); } break; case 'document': $text = $message['document']['filename'] ?? ''; $mediaUrl = $message['document']['url'] ?? $message['document']['id'] ?? null; if ($mediaUrl) { $this->queueMediaDownload($mediaUrl, $message['id'], 'document'); } break; case 'reaction': $emoji = $message['reaction']['emoji'] ?? ''; $reactionTo = $message['reaction']['message_id'] ?? ''; $text = json_encode(['emoji' => $emoji, 'message_id' => $reactionTo]); break; } return [ 'text' => $text, 'type' => $type, 'media_url' => $mediaUrl, 'timestamp' => $timestamp ]; } /** * Encolar descarga de media */ private function queueMediaDownload($mediaUrl, $messageId, $type) { $this->queue->push('media', [ 'media_url' => strpos($mediaUrl, 'http') === 0 ? $mediaUrl : null, 'media_id' => strpos($mediaUrl, 'http') !== 0 ? $mediaUrl : null, 'message_id' => $messageId, 'type' => $type, 'subdir' => date('Y/m') ], 1); // Prioridad normal } /** * Actualizar estados de mensajes (entregado, leído, etc.) */ private function updateMessageStatuses($statuses) { foreach ($statuses as $status) { try { $messageId = $status['id']; $newStatus = $status['status']; $this->db->update( 'conversations', ['status' => $newStatus], 'message_id = :message_id', ['message_id' => $messageId] ); } catch (Exception $e) { $this->logger->error("Failed to update message status", [ 'message_id' => $messageId ?? 'unknown', 'error' => $e->getMessage() ]); } } } /** * Guardar mensaje entrante en BD */ private function saveIncomingMessage($userId, $messageId, $text, $type, $mediaUrl, $timestamp) { try { $this->db->insert('conversations', [ 'user_id' => $userId, 'message_id' => $messageId, 'direction' => 'incoming', 'content' => $text, 'message_type' => $type, 'media_url' => $mediaUrl, 'status' => 'received', 'created_at' => $timestamp ? date('Y-m-d H:i:s', $timestamp) : date('Y-m-d H:i:s') ]); } catch (Exception $e) { $this->logger->error("Failed to save message", [ 'user_id' => $userId, 'message_id' => $messageId, 'error' => $e->getMessage() ]); } } /** * Log webhook de forma asíncrona (no bloqueante) */ private function logWebhookAsync($payload) { try { // Guardar solo últimos 1000 webhooks para no llenar BD $this->db->query("DELETE FROM webhook_logs WHERE id < (SELECT id FROM (SELECT id FROM webhook_logs ORDER BY id DESC LIMIT 1 OFFSET 1000) as t)"); $this->db->insert('webhook_logs', [ 'request_body' => $payload, 'response_body' => json_encode(['status' => 'queued']), 'status_code' => 200, 'created_at' => date('Y-m-d H:i:s') ]); } catch (Exception $e) { $this->logger->warning("Failed to log webhook", ['error' => $e->getMessage()]); } } private function getUserByPhone($phoneNumber) { return $this->db->fetch( "SELECT * FROM users WHERE phone_number = :phone", ['phone' => $phoneNumber] ); } private function createUser($phoneNumber) { $this->db->insert('users', [ 'phone_number' => $phoneNumber, 'name' => $phoneNumber, 'bot_enabled' => 1, 'created_at' => date('Y-m-d H:i:s') ]); return $this->db->lastInsertId(); } private function getUserById($userId) { return $this->db->fetch( "SELECT * FROM users WHERE id = :id", ['id' => $userId] ); } } // Ejecutar webhook try { $db = Database::getInstance(); $webhook = new OptimizedWebhook($db, $logger, $queue); $webhook->handleRequest(); } catch (Throwable $e) { $logger->critical("Webhook crashed", [ 'error' => $e->getMessage(), 'trace' => $e->getTraceAsString() ]); http_response_code(500); echo json_encode(['error' => 'Internal server error']); }