Files
whatsapp/api/webhook_optimized.php
T
2026-01-27 23:56:49 -05:00

424 lines
15 KiB
PHP

<?php
/**
* Webhook OPTIMIZADO para recibir mensajes de WhatsApp
*
* CAMBIOS CLAVE:
* - Responde inmediatamente (< 1 segundo) con 200 OK
* - Encola mensajes para procesamiento asíncrono
* - Usa Redis para evitar duplicados
* - Logging estructurado con Monolog
*
* Fecha: Enero 2026
*/
require_once __DIR__ . '/../vendor/autoload.php';
require_once __DIR__ . '/../config/config.php';
use WhatsApp\Queue\RedisQueue;
use Monolog\Logger;
use Monolog\Handler\StreamHandler;
use Monolog\Handler\RotatingFileHandler;
// Headers
header('Content-Type: application/json; charset=utf-8');
header('Access-Control-Allow-Origin: *');
header('Access-Control-Allow-Methods: GET, POST');
header('Access-Control-Allow-Headers: Content-Type');
// Configurar logger
$logger = new Logger('webhook');
$logger->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']);
}