up
This commit is contained in:
@@ -0,0 +1,360 @@
|
||||
<?php
|
||||
/**
|
||||
* Gestor de colas con Redis para procesamiento asíncrono
|
||||
* Implementa pattern Producer-Consumer para mensajes de WhatsApp
|
||||
*/
|
||||
|
||||
namespace WhatsApp\Queue;
|
||||
|
||||
use Predis\Client as RedisClient;
|
||||
use Monolog\Logger;
|
||||
|
||||
class RedisQueue {
|
||||
private $redis;
|
||||
private $logger;
|
||||
private $queuePrefix = 'whatsapp:queue:';
|
||||
|
||||
// Límites de rate limiting para WhatsApp Business API
|
||||
const RATE_LIMIT_WINDOW = 1; // segundos
|
||||
const MAX_MESSAGES_PER_SECOND = 80;
|
||||
const RATE_LIMIT_KEY = 'whatsapp:ratelimit:';
|
||||
|
||||
public function __construct(RedisClient $redis = null, Logger $logger = null) {
|
||||
// Conectar a Redis con configuración desde .env
|
||||
$this->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();
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user