#!/usr/bin/env php setFormatter($formatter); $logger->pushHandler($handler); // Handler para STDOUT en modo no-daemon if (!$isDaemon) { $consoleHandler = new StreamHandler('php://stdout', Logger::INFO); $consoleHandler->setFormatter($formatter); $logger->pushHandler($consoleHandler); } // Inicializar servicios $queue = new RedisQueue(null, $logger); $db = Database::getInstance(); $whatsappService = new WhatsAppService(); $botService = new BotService(); $logger->info("Worker started", [ 'queues' => $queues, 'daemon' => $isDaemon, 'pid' => getmypid() ]); // Manejadores de señales para shutdown graceful $shutdown = false; pcntl_async_signals(true); pcntl_signal(SIGTERM, function() use (&$shutdown, $logger) { $logger->info("Received SIGTERM, shutting down gracefully..."); $shutdown = true; }); pcntl_signal(SIGINT, function() use (&$shutdown, $logger) { $logger->info("Received SIGINT, shutting down gracefully..."); $shutdown = true; }); /** * Procesar un mensaje de la cola */ function processMessage($message, $queue, $botService, $whatsappService, $db, $logger) { $data = $message['data']; $attempts = $message['attempts']; try { $logger->debug("Processing message", [ 'queue' => $queue, 'attempts' => $attempts, 'data' => array_keys($data) ]); switch ($queue) { case 'messages': // Procesar mensaje de texto/interactivo if (!isset($data['user']) || !isset($data['messageText'])) { throw new Exception("Invalid message data: missing user or messageText"); } $user = $data['user']; $messageText = $data['messageText']; $messageType = $data['messageType'] ?? 'text'; // Procesar con BotService $botService->processMessage($user, $messageText, $messageType); $logger->info("Message processed successfully", [ 'user_id' => $user['id'], 'phone' => $user['phone_number'], 'type' => $messageType ]); break; case 'media': // Procesar descarga de media if (!isset($data['media_id']) && !isset($data['media_url'])) { throw new Exception("Invalid media data: missing media_id or media_url"); } $conversationId = $data['conversation_id'] ?? null; $mediaId = $data['media_id'] ?? null; $mediaUrl = $data['media_url'] ?? null; $subdir = $data['subdir'] ?? date('Y/m'); $mediaService = new MediaService(); if ($mediaId) { $result = $mediaService->fetchAndStoreFromGraph($mediaId, $subdir); } else { $result = $mediaService->fetchAndStoreFromUrl($mediaUrl, $subdir); } // Actualizar conversación con archivo local if ($result && $conversationId) { $update = []; if (!empty($result['local_file'])) { $update['local_file'] = $result['local_file']; } if (!empty($result['local_thumb'])) { $update['local_thumb'] = $result['local_thumb']; } if (!empty($update)) { $db->update('conversations', $update, 'id = :id', ['id' => $conversationId]); } } $logger->info("Media processed successfully", [ 'conversation_id' => $conversationId, 'media_id' => $mediaId, 'local_file' => $result['local_file'] ?? null ]); break; case 'notifications': // Procesar notificación (enviar push, email, etc.) if (!isset($data['type']) || !isset($data['message'])) { throw new Exception("Invalid notification data"); } $type = $data['type']; $notificationMessage = $data['message']; $userId = $data['user_id'] ?? null; // Aquí puedes integrar servicios de notificación // Por ahora solo lo registramos $logger->info("Notification processed", [ 'type' => $type, 'user_id' => $userId, 'message' => substr($notificationMessage, 0, 100) ]); break; default: $logger->warning("Unknown queue type", ['queue' => $queue]); } return true; } catch (Exception $e) { $logger->error("Failed to process message", [ 'queue' => $queue, 'error' => $e->getMessage(), 'trace' => $e->getTraceAsString(), 'attempts' => $attempts ]); return false; } } /** * Loop principal del worker */ $processedCount = 0; $errorCount = 0; $maxErrors = 10; // Detener si hay muchos errores consecutivos do { try { // Procesar mensajes delayed cada 10 iteraciones if ($processedCount % 10 === 0) { foreach ($queues as $q) { $queue->processDelayed($q); } } // Esperar mensaje de cualquier cola (timeout 5 segundos) $message = $queue->pop($queues, 5); if (!$message) { // No hay mensajes, continuar esperando if (!$isDaemon) { // En modo no-daemon, salir si no hay más mensajes $logger->info("No more messages, exiting"); break; } continue; } // Procesar mensaje $success = processMessage( $message, $message['queue'], $botService, $whatsappService, $db, $logger ); if ($success) { $processedCount++; $errorCount = 0; // Reset error counter } else { $errorCount++; // Re-encolar con retry if (!$queue->retry($message['queue'], $message, 3)) { $logger->critical("Failed to retry message", [ 'queue' => $message['queue'] ]); } // Si hay muchos errores consecutivos, detener worker if ($errorCount >= $maxErrors) { $logger->critical("Too many consecutive errors, stopping worker", [ 'error_count' => $errorCount ]); break; } } // Limitar memoria (reiniciar si supera 128MB) $memoryUsage = memory_get_usage(true) / 1024 / 1024; if ($memoryUsage > 128) { $logger->warning("Memory limit reached, restarting worker", [ 'memory_mb' => round($memoryUsage, 2) ]); break; } } catch (Exception $e) { $logger->error("Worker error", [ 'error' => $e->getMessage(), 'trace' => $e->getTraceAsString() ]); $errorCount++; if ($errorCount >= $maxErrors) { break; } // Esperar un poco antes de reintentar sleep(5); } } while ($isDaemon && !$shutdown); // Cleanup $queue->disconnect(); $logger->info("Worker stopped", [ 'processed' => $processedCount, 'errors' => $errorCount ]); exit(0);