redis = $redis ?? new RedisClient([ 'scheme' => getenv('REDIS_SCHEME') ?: 'tcp', 'host' => getenv('REDIS_HOST') ?: '127.0.0.1', 'port' => getenv('REDIS_PORT') ?: 6379, 'password' => getenv('REDIS_PASSWORD') ?: null, 'database' => getenv('REDIS_DB') ?: 0, ]); $this->logger = $logger; } /** * Agregar mensaje a la cola para procesamiento asíncrono * * @param string $queueName Nombre de la cola (ej: 'messages', 'media', 'notifications') * @param array $data Datos del mensaje a procesar * @param int $priority Prioridad (0=alta, 1=normal, 2=baja) * @return bool */ public function push(string $queueName, array $data, int $priority = 1): bool { try { $payload = json_encode([ 'data' => $data, 'priority' => $priority, 'queued_at' => time(), 'attempts' => 0 ]); // Usar LPUSH para agregar al inicio (FIFO con BRPOP) $key = $this->queuePrefix . $queueName; $result = $this->redis->lpush($key, [$payload]); if ($this->logger) { $this->logger->info("Message pushed to queue", [ 'queue' => $queueName, 'priority' => $priority, 'data_keys' => array_keys($data) ]); } return $result > 0; } catch (\Exception $e) { if ($this->logger) { $this->logger->error("Failed to push to queue", [ 'queue' => $queueName, 'error' => $e->getMessage() ]); } return false; } } /** * Obtener mensaje de la cola (bloqueante) * * @param string|array $queueName Nombre de cola(s) a escuchar * @param int $timeout Timeout en segundos (0 = infinito) * @return array|null [queue_name, payload] o null si timeout */ public function pop($queueName, int $timeout = 5): ?array { try { $queues = is_array($queueName) ? $queueName : [$queueName]; $keys = array_map(fn($q) => $this->queuePrefix . $q, $queues); // BRPOP espera hasta que haya un elemento o timeout $result = $this->redis->brpop($keys, $timeout); if (!$result) { return null; } // $result = [queue_key, payload_json] $queueKey = $result[0]; $payload = json_decode($result[1], true); // Extraer nombre de cola sin prefijo $queue = str_replace($this->queuePrefix, '', $queueKey); if ($this->logger) { $this->logger->debug("Message popped from queue", [ 'queue' => $queue, 'attempts' => $payload['attempts'] ?? 0 ]); } return [ 'queue' => $queue, 'data' => $payload['data'] ?? [], 'priority' => $payload['priority'] ?? 1, 'queued_at' => $payload['queued_at'] ?? time(), 'attempts' => $payload['attempts'] ?? 0 ]; } catch (\Exception $e) { if ($this->logger) { $this->logger->error("Failed to pop from queue", [ 'queue' => $queueName, 'error' => $e->getMessage() ]); } return null; } } /** * Re-encolar mensaje fallido con backoff exponencial * * @param string $queueName * @param array $message Mensaje original con metadata * @param int $maxAttempts Intentos máximos antes de mover a DLQ * @return bool */ public function retry(string $queueName, array $message, int $maxAttempts = 3): bool { $attempts = ($message['attempts'] ?? 0) + 1; if ($attempts >= $maxAttempts) { // Mover a Dead Letter Queue return $this->moveToDLQ($queueName, $message); } // Incrementar contador de intentos $message['attempts'] = $attempts; $message['last_attempt'] = time(); // Backoff exponencial: 2^attempts segundos $delay = pow(2, $attempts); try { $payload = json_encode($message); $delayedKey = $this->queuePrefix . $queueName . ':delayed'; // Usar ZADD con score = timestamp futuro $executeAt = time() + $delay; $this->redis->zadd($delayedKey, [$payload => $executeAt]); if ($this->logger) { $this->logger->warning("Message scheduled for retry", [ 'queue' => $queueName, 'attempt' => $attempts, 'delay' => $delay, 'execute_at' => date('Y-m-d H:i:s', $executeAt) ]); } return true; } catch (\Exception $e) { if ($this->logger) { $this->logger->error("Failed to schedule retry", [ 'queue' => $queueName, 'error' => $e->getMessage() ]); } return false; } } /** * Mover mensaje a Dead Letter Queue (mensajes que fallaron definitivamente) */ private function moveToDLQ(string $queueName, array $message): bool { try { $dlqKey = $this->queuePrefix . 'dlq:' . $queueName; $payload = json_encode([ 'original_message' => $message, 'failed_at' => time(), 'attempts' => $message['attempts'] ?? 0 ]); $this->redis->lpush($dlqKey, [$payload]); if ($this->logger) { $this->logger->error("Message moved to DLQ", [ 'queue' => $queueName, 'attempts' => $message['attempts'] ?? 0 ]); } return true; } catch (\Exception $e) { if ($this->logger) { $this->logger->critical("Failed to move to DLQ", [ 'queue' => $queueName, 'error' => $e->getMessage() ]); } return false; } } /** * Procesar mensajes delayed (ejecutar cron cada minuto) */ public function processDelayed(string $queueName): int { try { $delayedKey = $this->queuePrefix . $queueName . ':delayed'; $now = time(); // Obtener mensajes cuyo score (timestamp) <= ahora $messages = $this->redis->zrangebyscore($delayedKey, '-inf', $now); $processed = 0; foreach ($messages as $payload) { $message = json_decode($payload, true); // Mover de delayed a queue normal if ($this->push($queueName, $message['data'], $message['priority'] ?? 1)) { $this->redis->zrem($delayedKey, $payload); $processed++; } } if ($processed > 0 && $this->logger) { $this->logger->info("Processed delayed messages", [ 'queue' => $queueName, 'count' => $processed ]); } return $processed; } catch (\Exception $e) { if ($this->logger) { $this->logger->error("Failed to process delayed", [ 'queue' => $queueName, 'error' => $e->getMessage() ]); } return 0; } } /** * Verificar rate limit para WhatsApp API * * @param string $identifier Identificador único (ej: phone_number_id) * @return bool true si puede enviar, false si excede límite */ public function checkRateLimit(string $identifier): bool { try { $key = self::RATE_LIMIT_KEY . $identifier; $count = $this->redis->incr($key); // Establecer expiración solo en el primer incremento if ($count == 1) { $this->redis->expire($key, self::RATE_LIMIT_WINDOW); } if ($count > self::MAX_MESSAGES_PER_SECOND) { if ($this->logger) { $this->logger->warning("Rate limit exceeded", [ 'identifier' => $identifier, 'count' => $count, 'limit' => self::MAX_MESSAGES_PER_SECOND ]); } return false; } return true; } catch (\Exception $e) { if ($this->logger) { $this->logger->error("Rate limit check failed", [ 'error' => $e->getMessage() ]); } // En caso de error, permitir para no bloquear return true; } } /** * Obtener estadísticas de las colas */ public function getStats(array $queueNames = ['messages', 'media', 'notifications']): array { $stats = []; foreach ($queueNames as $queueName) { $key = $this->queuePrefix . $queueName; $delayedKey = $key . ':delayed'; $dlqKey = $this->queuePrefix . 'dlq:' . $queueName; $stats[$queueName] = [ 'pending' => $this->redis->llen($key), 'delayed' => $this->redis->zcard($delayedKey), 'failed' => $this->redis->llen($dlqKey) ]; } return $stats; } /** * Limpiar colas (para testing/desarrollo) */ public function flush(string $queueName = null): bool { try { if ($queueName) { $keys = [ $this->queuePrefix . $queueName, $this->queuePrefix . $queueName . ':delayed', $this->queuePrefix . 'dlq:' . $queueName ]; foreach ($keys as $key) { $this->redis->del([$key]); } } else { // Flush todas las colas con el prefijo $pattern = $this->queuePrefix . '*'; $keys = $this->redis->keys($pattern); if (!empty($keys)) { $this->redis->del($keys); } } if ($this->logger) { $this->logger->info("Queue flushed", ['queue' => $queueName ?? 'all']); } return true; } catch (\Exception $e) { if ($this->logger) { $this->logger->error("Failed to flush queue", [ 'queue' => $queueName, 'error' => $e->getMessage() ]); } return false; } } /** * Cerrar conexión Redis */ public function disconnect(): void { if ($this->redis) { $this->redis->disconnect(); } } }