Files
projet_gestion_commande/backend/gestion/db/redis_queue_management.go
T

748 lines
21 KiB
Go

package db
import (
"encoding/json"
"fmt"
"gestion/models"
"log"
"strconv"
"time"
"github.com/redis/go-redis/v9"
)
// ============================================
// 🆕 GESTION AUTOMATIQUE DU STATUT BUSY
// ============================================
// UpdateDeliverymanStatusBasedOnQueue met à jour automatiquement le statut
// ✅ Status = "busy" si queue >= 10
// ✅ Status = "available" si queue < 10
func (d *Database) UpdateDeliverymanStatusBasedOnQueue(deliveryman string) error {
queueKey := fmt.Sprintf("queue:deliveryman:%s", deliveryman)
queueSize, err := Redis.ZCard(RedisCtx, queueKey).Result()
if err != nil {
return fmt.Errorf("erreur récupération taille queue: %w", err)
}
// Récupérer le statut actuel
statusKey := fmt.Sprintf("delivery:status:%s", deliveryman)
data, err := Redis.Get(RedisCtx, statusKey).Result()
if err != nil {
log.Printf("⚠️ [STATUS] Livreur %s n'a pas de statut Redis", deliveryman)
return nil
}
var status models.DeliveryPersonStatus
if err := json.Unmarshal([]byte(data), &status); err != nil {
return err
}
// Déterminer le nouveau statut
var newStatus string
var currentCommand int
if queueSize >= MAX_COMMANDS_PER_DELIVERYMAN {
// 🔴 BUSY car queue pleine (10 commandes ou plus)
newStatus = "busy"
currentCommand = 0
log.Printf("🔴 [STATUS] %s -> BUSY (queue pleine: %d/10)", deliveryman, queueSize)
} else {
// 🟢 AVAILABLE tant que queue < 10
// Exception: si le livreur est en train de livrer (delivering), on garde ce statut
if status.Status == "delivering" && status.CurrentCommand > 0 {
newStatus = "delivering"
currentCommand = status.CurrentCommand
log.Printf("🟡 [STATUS] %s -> DELIVERING (queue: %d/10, livraison en cours: cmd %d)",
deliveryman, queueSize, currentCommand)
} else {
newStatus = "available"
currentCommand = 0
log.Printf("🟢 [STATUS] %s -> AVAILABLE (queue: %d/10)", deliveryman, queueSize)
}
}
// Mettre à jour le statut
return d.SetDeliveryPersonStatus(deliveryman, newStatus, currentCommand)
}
// CanDeliverymanAcceptCommands vérifie si un livreur peut accepter de nouvelles commandes
// ✅ Retourne false si: status=busy ET queue>=10, ou status=offline
func (d *Database) CanDeliverymanAcceptCommands(deliveryman string) bool {
// 1. Vérifier le statut Redis
statusKey := fmt.Sprintf("delivery:status:%s", deliveryman)
data, err := Redis.Get(RedisCtx, statusKey).Result()
if err != nil {
log.Printf("⚠️ [CHECK] Livreur %s sans statut Redis", deliveryman)
return true // Fallback: autoriser si pas de statut
}
var status models.DeliveryPersonStatus
if err := json.Unmarshal([]byte(data), &status); err != nil {
return true
}
// 2. Si offline, refuser
if status.Status == "offline" {
log.Printf("⚫ [CHECK] %s REFUSÉ: offline", deliveryman)
return false
}
// 3. Vérifier la taille de la queue
queueKey := fmt.Sprintf("queue:deliveryman:%s", deliveryman)
queueSize, _ := Redis.ZCard(RedisCtx, queueKey).Result()
// 4. Si queue >= 10, refuser
if queueSize >= MAX_COMMANDS_PER_DELIVERYMAN {
log.Printf("🔴 [CHECK] %s REFUSÉ: queue pleine (%d/10)", deliveryman, queueSize)
return false
}
log.Printf("🟢 [CHECK] %s AUTORISÉ (%d/10)", deliveryman, queueSize)
return true
}
// GetAvailableDeliveryPersonsForAssignment récupère UNIQUEMENT les livreurs pouvant accepter
func (d *Database) GetAvailableDeliveryPersonsForAssignment() ([]models.DeliveryPersonStatus, error) {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return nil, err
}
var available []models.DeliveryPersonStatus
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
if err := json.Unmarshal([]byte(data), &status); err != nil {
continue
}
// ✅ Vérifier si le livreur peut accepter des commandes
if d.CanDeliverymanAcceptCommands(status.Username) {
available = append(available, status)
}
}
log.Printf("📊 [AVAILABLE] %d livreur(s) disponible(s) pour assignation", len(available))
return available, nil
}
// ============================================
// 🔄 FONCTIONS MODIFIÉES AVEC AUTO-STATUS
// ============================================
// AddToDeliverymanQueue - VERSION MISE À JOUR avec auto-update du statut
func (d *Database) AddToDeliverymanQueue(deliveryman string, queueItem models.CommandQueue) error {
data, err := json.Marshal(queueItem)
if err != nil {
return fmt.Errorf("erreur serialisation: %w", err)
}
// Sauvegarder les données de la commande
commandKey := fmt.Sprintf("queue:pending:%d", queueItem.CommandID)
Redis.Set(RedisCtx, commandKey, data, 24*time.Hour)
// Ajouter à la queue sorted set du livreur (score = timestamp)
queueKey := fmt.Sprintf("queue:deliveryman:%s", deliveryman)
score := float64(time.Now().Unix())
err = Redis.ZAdd(RedisCtx, queueKey, redis.Z{
Score: score,
Member: queueItem.CommandID,
}).Err()
if err != nil {
return fmt.Errorf("erreur ajout queue Redis: %w", err)
}
// Incrémenter le compteur de commandes en attente pour ce livreur
counterKey := fmt.Sprintf("queue:deliveryman:%s:count", deliveryman)
Redis.Incr(RedisCtx, counterKey)
log.Printf("✅ Commande %d ajoutée à la queue de %s", queueItem.CommandID, deliveryman)
// ✅ NOUVEAU: Mettre à jour automatiquement le statut
go d.UpdateDeliverymanStatusBasedOnQueue(deliveryman)
return nil
}
// AddCommandToQueue ajoute une commande à la file d'attente Redis (version simple)
func (d *Database) AddCommandToQueue(commandID int) error {
if err := d.ValidateCommandBeforeQueue(commandID); err != nil {
log.Printf("❌ [QUEUE] Commande %d REFUSÉE: %v", commandID, err)
return fmt.Errorf("validation échouée: %w", err)
}
command, err := d.GetCommandByID(commandID)
if err != nil {
return fmt.Errorf("commande introuvable: %w", err)
}
// Gestion des coordonnées (NULL safe)
var lat, lng float64
if command["dest_latitude"] != nil {
if latVal, ok := command["dest_latitude"].(float64); ok {
lat = latVal
}
}
if command["dest_longitude"] != nil {
if lngVal, ok := command["dest_longitude"].(float64); ok {
lng = lngVal
}
}
// Récupérer le total_prix avec gestion de type
var totalPrice float64
if tp, ok := command["total_prix"].(float64); ok {
totalPrice = tp
}
// Récupérer l'adresse
var address string
if addr, ok := command["delivery_address"].(string); ok {
address = addr
}
queueItem := models.CommandQueue{
CommandID: commandID,
Username: command["username"].(string),
TotalPrice: totalPrice,
Address: address,
Lat: lat,
Lng: lng,
CreatedAt: time.Now(),
EstimatedETA: 0,
}
// Ajouter à la queue générale
return d.AddToGeneralQueue(queueItem)
}
// AddCommandToSmartQueue - Ajoute une commande avec attribution au livreur le moins chargé
func (d *Database) AddCommandToSmartQueue(commandID int, address string) error {
if err := d.ValidateCommandBeforeQueue(commandID); err != nil {
log.Printf("❌ [QUEUE] Commande %d REFUSÉE: %v", commandID, err)
return fmt.Errorf("validation échouée: %w", err)
}
command, err := d.GetCommandByID(commandID)
if err != nil {
return fmt.Errorf("commande introuvable: %w", err)
}
// Gestion des coordonnées (NULL safe)
var lat, lng float64
if command["dest_latitude"] != nil {
if latVal, ok := command["dest_latitude"].(float64); ok {
lat = latVal
}
}
if command["dest_longitude"] != nil {
if lngVal, ok := command["dest_longitude"].(float64); ok {
lng = lngVal
}
}
// Récupérer le total_prix avec gestion de type
var totalPrice float64
if tp, ok := command["total_prix"].(float64); ok {
totalPrice = tp
}
queueItem := models.CommandQueue{
CommandID: commandID,
Username: command["username"].(string),
TotalPrice: totalPrice,
Address: address,
Lat: lat,
Lng: lng,
CreatedAt: time.Now(),
EstimatedETA: 0,
}
// ✅ MODIFIÉ: Utiliser FindLeastLoadedDeliveryman qui respecte maintenant le statut
assignedDeliveryman, err := d.FindLeastLoadedDeliveryman()
if err != nil {
// Aucun livreur trouvé, ajouter à la queue générale
log.Printf("⚠️ Aucun livreur trouvé, ajout à la queue générale")
return d.AddToGeneralQueue(queueItem)
}
// Ajouter la commande à la queue spécifique du livreur (auto-update du statut)
err = d.AddToDeliverymanQueue(assignedDeliveryman, queueItem)
if err != nil {
return fmt.Errorf("erreur ajout à la queue du livreur: %w", err)
}
log.Printf("📋 Commande %d assignée à la queue de %s", commandID, assignedDeliveryman)
// Publier l'événement
d.PublishCommandEvent(commandID, "queued",
fmt.Sprintf("En attente dans la queue de %s", assignedDeliveryman))
return nil
}
// AddToGeneralQueue ajoute une commande à la queue générale (fallback)
func (d *Database) AddToGeneralQueue(queueItem models.CommandQueue) error {
data, err := json.Marshal(queueItem)
if err != nil {
return fmt.Errorf("erreur serialisation: %w", err)
}
score := float64(time.Now().Unix())
key := fmt.Sprintf("queue:pending:%d", queueItem.CommandID)
pipe := Redis.Pipeline()
pipe.Set(RedisCtx, key, data, 24*time.Hour)
pipe.ZAdd(RedisCtx, "queue:pending:sorted", redis.Z{
Score: score,
Member: queueItem.CommandID,
})
_, err = pipe.Exec(RedisCtx)
if err != nil {
return fmt.Errorf("erreur ajout queue générale: %w", err)
}
log.Printf("📥 Commande %d ajoutée à la queue générale", queueItem.CommandID)
return nil
}
// RemoveCommandFromQueue - VERSION AMÉLIORÉE avec auto-update du statut
func (d *Database) RemoveCommandFromQueue(commandID int) error {
key := fmt.Sprintf("queue:pending:%d", commandID)
commandIDStr := strconv.Itoa(commandID)
pipe := Redis.Pipeline()
pipe.Del(RedisCtx, key)
pipe.ZRem(RedisCtx, "queue:pending:sorted", commandIDStr)
pipe.ZRem(RedisCtx, "queue:priority:sorted", commandIDStr)
// Trouver et retirer de la queue du livreur
keys, _ := Redis.Keys(RedisCtx, "queue:deliveryman:*").Result()
var affectedDeliveryman string
for _, queueKey := range keys {
if len(queueKey) > 6 && queueKey[len(queueKey)-6:] == ":count" {
continue
}
// Vérifier si la commande est dans cette queue
_, err := Redis.ZRank(RedisCtx, queueKey, commandIDStr).Result()
if err == nil {
// Commande trouvée dans cette queue
affectedDeliveryman = queueKey[len("queue:deliveryman:"):]
pipe.ZRem(RedisCtx, queueKey, commandIDStr)
// Décrémenter le compteur
counterKey := fmt.Sprintf("queue:deliveryman:%s:count", affectedDeliveryman)
pipe.Decr(RedisCtx, counterKey)
break
}
}
_, err := pipe.Exec(RedisCtx)
if err != nil {
return fmt.Errorf("erreur suppression: %w", err)
}
log.Printf("✅ Commande %d retirée de la file", commandID)
// ✅ NOUVEAU: Mettre à jour le statut si un livreur était affecté
if affectedDeliveryman != "" {
go d.UpdateDeliverymanStatusBasedOnQueue(affectedDeliveryman)
}
return nil
}
func (d *Database) GetNextCommandInQueue() (*models.CommandQueue, error) {
normalResults, err := Redis.ZRangeWithScores(RedisCtx, "queue:pending:sorted", 0, 0).Result()
if err != nil || len(normalResults) == 0 {
return nil, fmt.Errorf("aucune commande en attente")
}
commandID := extractCommandID(normalResults[0].Member)
if commandID <= 0 {
return nil, fmt.Errorf("ID commande invalide")
}
key := fmt.Sprintf("queue:pending:%d", commandID)
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
return nil, fmt.Errorf("commande introuvable: %w", err)
}
var queue models.CommandQueue
if err := json.Unmarshal([]byte(data), &queue); err != nil {
return nil, fmt.Errorf("erreur désérialisation: %w", err)
}
return &queue, nil
}
// GetLastCommandInQueue récupère la dernière commande dans la queue d'un livreur
func (d *Database) GetLastCommandInQueue(deliveryman string) (*models.CommandQueue, error) {
queueKey := fmt.Sprintf("queue:deliveryman:%s", deliveryman)
// Récupérer la dernière commande (index -1)
commandIDs, err := Redis.ZRange(RedisCtx, queueKey, -1, -1).Result()
if err != nil || len(commandIDs) == 0 {
return nil, fmt.Errorf("queue vide")
}
commandID := extractCommandID(commandIDs[0])
if commandID <= 0 {
return nil, fmt.Errorf("ID invalide")
}
commandKey := fmt.Sprintf("queue:pending:%d", commandID)
data, err := Redis.Get(RedisCtx, commandKey).Result()
if err != nil {
return nil, err
}
var queueItem models.CommandQueue
if err := json.Unmarshal([]byte(data), &queueItem); err != nil {
return nil, err
}
return &queueItem, nil
}
// GetCommandQueuePosition récupère la position d'une commande dans la queue
func (d *Database) GetCommandQueuePosition(commandID int) (int, error) {
commandIDStr := strconv.Itoa(commandID)
// Chercher d'abord dans les queues des livreurs
keys, _ := Redis.Keys(RedisCtx, "queue:deliveryman:*").Result()
for _, queueKey := range keys {
// Éviter les clés de compteur
if len(queueKey) > 6 && queueKey[len(queueKey)-6:] == ":count" {
continue
}
rank, err := Redis.ZRank(RedisCtx, queueKey, commandIDStr).Result()
if err == nil {
return int(rank) + 1, nil
}
}
// Chercher dans la queue générale
rank, err := Redis.ZRank(RedisCtx, "queue:pending:sorted", commandIDStr).Result()
if err == nil {
return int(rank) + 1, nil
}
return 0, fmt.Errorf("commande non trouvée dans les queues")
}
func (d *Database) ClearDeliverymanQueue(deliveryman string) error {
queueKey := fmt.Sprintf("queue:deliveryman:%s", deliveryman)
// Récupérer toutes les commandes
commandIDs, _ := Redis.ZRange(RedisCtx, queueKey, 0, -1).Result()
// Redistribuer chaque commande
for _, cmdIDStr := range commandIDs {
commandID := extractCommandID(cmdIDStr)
if commandID <= 0 {
continue
}
// Récupérer les données de la commande
commandKey := fmt.Sprintf("queue:pending:%d", commandID)
data, err := Redis.Get(RedisCtx, commandKey).Result()
if err != nil {
continue
}
var queueItem models.CommandQueue
if err := json.Unmarshal([]byte(data), &queueItem); err != nil {
continue
}
// Trouver un nouveau livreur
newDeliveryman, err := d.FindLeastLoadedDeliveryman()
if err != nil {
// Fallback: queue générale
d.AddToGeneralQueue(queueItem)
continue
}
// Réassigner à un autre livreur
if newDeliveryman != deliveryman {
d.AddToDeliverymanQueue(newDeliveryman, queueItem)
log.Printf("🔄 Commande %d réassignée de %s à %s",
commandID, deliveryman, newDeliveryman)
}
}
// Vider la queue
Redis.Del(RedisCtx, queueKey)
Redis.Del(RedisCtx, fmt.Sprintf("queue:deliveryman:%s:count", deliveryman))
log.Printf("🗑️ Queue de %s vidée et redistribuée", deliveryman)
// ✅ NOUVEAU: Mettre à jour le statut
go d.UpdateDeliverymanStatusBasedOnQueue(deliveryman)
return nil
}
func (d *Database) SetDeliveryPersonStatus(username, status string, commandID int) error {
key := fmt.Sprintf("delivery:status:%s", username)
statusData := models.DeliveryPersonStatus{
Username: username,
Status: status,
CurrentCommand: commandID,
LastUpdate: time.Now(),
}
data, _ := json.Marshal(statusData)
err := Redis.Set(RedisCtx, key, data, 24*time.Hour).Err()
if err == nil {
log.Printf("✅ Statut livreur %s: %s", username, status)
}
return err
}
// GetAvailableDeliveryPersonsRedis récupère les livreurs disponibles (LEGACY)
func (d *Database) GetAvailableDeliveryPersonsRedis() ([]models.DeliveryPersonStatus, error) {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return nil, err
}
var available []models.DeliveryPersonStatus
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
json.Unmarshal([]byte(data), &status)
if status.Status == "available" {
available = append(available, status)
}
}
return available, nil
}
// GetAllActiveDeliveryPersons - VERSION MISE À JOUR avec vérification capacité
func (d *Database) GetAllActiveDeliveryPersons() ([]models.DeliveryPersonStatus, error) {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return nil, err
}
var active []models.DeliveryPersonStatus
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
json.Unmarshal([]byte(data), &status)
if status.Status != "offline" {
// ✅ Vérifier si le livreur peut accepter des commandes
if d.CanDeliverymanAcceptCommands(status.Username) {
active = append(active, status)
}
}
}
return active, nil
}
// CountActiveDeliverymen compte le nombre de livreurs actifs (non offline)
func (d *Database) CountActiveDeliverymen() (int, error) {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return 0, err
}
count := 0
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
json.Unmarshal([]byte(data), &status)
if status.Status != "offline" {
count++
}
}
return count, nil
}
func (d *Database) GetSingleActiveDeliveryman() (string, error) {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return "", err
}
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
json.Unmarshal([]byte(data), &status)
if status.Status != "offline" {
return status.Username, nil
}
}
return "", fmt.Errorf("aucun livreur actif trouvé")
}
// GetAllActiveDeliverymenUsernames retourne les usernames de tous les livreurs actifs
func (d *Database) GetAllActiveDeliverymenUsernames() ([]string, error) {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return nil, err
}
var activeUsernames []string
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
json.Unmarshal([]byte(data), &status)
if status.Status != "offline" {
activeUsernames = append(activeUsernames, status.Username)
}
}
return activeUsernames, nil
}
// ============================================
// 🔄 SYNCHRONISATION DES STATUTS
// ============================================
// SyncAllDeliverymanStatuses synchronise tous les statuts (à appeler au démarrage)
func (d *Database) SyncAllDeliverymanStatuses() error {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return err
}
log.Println("🔄 [SYNC] Synchronisation des statuts livreurs...")
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
if err := json.Unmarshal([]byte(data), &status); err != nil {
continue
}
// Mettre à jour le statut basé sur la queue
d.UpdateDeliverymanStatusBasedOnQueue(status.Username)
}
log.Println("✅ [SYNC] Synchronisation terminée")
return nil
}
// ============================================
// 📊 RAPPORT DE CAPACITÉ
// ============================================
// GetDeliverymanCapacityReport génère un rapport détaillé
func (d *Database) GetDeliverymanCapacityReport() (map[string]interface{}, error) {
keys, err := Redis.Keys(RedisCtx, "delivery:status:*").Result()
if err != nil {
return nil, err
}
report := map[string]interface{}{
"total_deliverymen": 0,
"available": 0,
"busy_full": 0, // BUSY car queue pleine
"busy_delivering": 0, // BUSY car en livraison
"offline": 0,
"details": []map[string]interface{}{},
}
for _, key := range keys {
data, err := Redis.Get(RedisCtx, key).Result()
if err != nil {
continue
}
var status models.DeliveryPersonStatus
if err := json.Unmarshal([]byte(data), &status); err != nil {
continue
}
queueKey := fmt.Sprintf("queue:deliveryman:%s", status.Username)
queueSize, _ := Redis.ZCard(RedisCtx, queueKey).Result()
canAccept := d.CanDeliverymanAcceptCommands(status.Username)
detail := map[string]interface{}{
"username": status.Username,
"status": status.Status,
"queue_size": queueSize,
"capacity": fmt.Sprintf("%d/10", queueSize),
"can_accept": canAccept,
"current_order": status.CurrentCommand,
}
report["total_deliverymen"] = report["total_deliverymen"].(int) + 1
if status.Status == "offline" {
report["offline"] = report["offline"].(int) + 1
} else if status.Status == "busy" {
if queueSize >= MAX_COMMANDS_PER_DELIVERYMAN {
report["busy_full"] = report["busy_full"].(int) + 1
} else {
report["busy_delivering"] = report["busy_delivering"].(int) + 1
}
} else if canAccept {
report["available"] = report["available"].(int) + 1
}
report["details"] = append(report["details"].([]map[string]interface{}), detail)
}
return report, nil
}