341 lines
9.7 KiB
Go
341 lines
9.7 KiB
Go
package db
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"gestion/models"
|
|
"log"
|
|
"strconv"
|
|
"time"
|
|
|
|
"github.com/redis/go-redis/v9"
|
|
)
|
|
|
|
// UpdateDeliverymanStatusBasedOnQueue met à jour automatiquement le statut
|
|
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)
|
|
}
|
|
|
|
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
|
|
}
|
|
|
|
var newStatus string
|
|
var currentCommand int
|
|
|
|
if queueSize >= MAX_COMMANDS_PER_DELIVERYMAN {
|
|
newStatus = "busy"
|
|
currentCommand = 0
|
|
log.Printf("🔴 [STATUS] %s -> BUSY (queue pleine: %d/10)", deliveryman, queueSize)
|
|
} else {
|
|
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)
|
|
}
|
|
}
|
|
|
|
return d.SetDeliveryPersonStatus(deliveryman, newStatus, currentCommand)
|
|
}
|
|
|
|
// CanDeliverymanAcceptCommands vérifie si un livreur peut accepter de nouvelles commandes
|
|
func (d *Database) CanDeliverymanAcceptCommands(deliveryman string) bool {
|
|
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
|
|
}
|
|
|
|
var status models.DeliveryPersonStatus
|
|
if err := json.Unmarshal([]byte(data), &status); err != nil {
|
|
return true
|
|
}
|
|
|
|
if status.Status == "offline" {
|
|
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
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
|
|
commandKey := fmt.Sprintf("queue:pending:%d", queueItem.CommandID)
|
|
Redis.Set(RedisCtx, commandKey, data, 24*time.Hour)
|
|
|
|
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)
|
|
}
|
|
|
|
counterKey := fmt.Sprintf("queue:deliveryman:%s:count", deliveryman)
|
|
Redis.Incr(RedisCtx, counterKey)
|
|
|
|
go d.UpdateDeliverymanStatusBasedOnQueue(deliveryman)
|
|
|
|
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, _ := scanRedisKeys("queue:deliveryman:*")
|
|
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
|
|
}
|
|
|
|
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) {
|
|
var available []models.DeliveryPersonStatus
|
|
err := d.iterDeliveryStatuses(func(s models.DeliveryPersonStatus) {
|
|
if s.Status == "available" {
|
|
available = append(available, s)
|
|
}
|
|
})
|
|
return available, err
|
|
}
|
|
|
|
// GetAllActiveDeliveryPersons retourne les livreurs actifs pouvant accepter des commandes
|
|
func (d *Database) GetAllActiveDeliveryPersons() ([]models.DeliveryPersonStatus, error) {
|
|
var active []models.DeliveryPersonStatus
|
|
err := d.iterDeliveryStatuses(func(s models.DeliveryPersonStatus) {
|
|
if s.Status != "offline" && d.CanDeliverymanAcceptCommands(s.Username) {
|
|
active = append(active, s)
|
|
}
|
|
})
|
|
return active, err
|
|
}
|
|
|
|
// CountActiveDeliverymen compte le nombre de livreurs actifs (non offline)
|
|
func (d *Database) CountActiveDeliverymen() (int, error) {
|
|
count := 0
|
|
err := d.iterDeliveryStatuses(func(s models.DeliveryPersonStatus) {
|
|
if s.Status != "offline" {
|
|
count++
|
|
}
|
|
})
|
|
return count, err
|
|
}
|
|
|
|
func (d *Database) GetSingleActiveDeliveryman() (string, error) {
|
|
found := ""
|
|
d.iterDeliveryStatuses(func(s models.DeliveryPersonStatus) { //nolint
|
|
if found == "" && s.Status != "offline" {
|
|
found = s.Username
|
|
}
|
|
})
|
|
if found == "" {
|
|
return "", fmt.Errorf("aucun livreur actif trouvé")
|
|
}
|
|
return found, nil
|
|
}
|
|
|
|
// GetAllActiveDeliverymenUsernames retourne les usernames de tous les livreurs actifs
|
|
func (d *Database) GetAllActiveDeliverymenUsernames() ([]string, error) {
|
|
var usernames []string
|
|
err := d.iterDeliveryStatuses(func(s models.DeliveryPersonStatus) {
|
|
if s.Status != "offline" {
|
|
usernames = append(usernames, s.Username)
|
|
}
|
|
})
|
|
return usernames, err
|
|
}
|
|
|
|
// SyncAllDeliverymanStatuses synchronise tous les statuts (à appeler au démarrage)
|
|
func (d *Database) SyncAllDeliverymanStatuses() error {
|
|
log.Println("🔄 [SYNC] Synchronisation des statuts livreurs...")
|
|
return d.iterDeliveryStatuses(func(s models.DeliveryPersonStatus) {
|
|
d.UpdateDeliverymanStatusBasedOnQueue(s.Username)
|
|
})
|
|
}
|
|
|
|
// iterDeliveryStatuses itère sur tous les statuts Redis des livreurs et appelle fn pour chacun.
|
|
func (d *Database) iterDeliveryStatuses(fn func(models.DeliveryPersonStatus)) error {
|
|
keys, err := scanRedisKeys("delivery:status:*")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
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
|
|
}
|
|
fn(status)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// scanRedisKeys remplace KEYS * par SCAN pour ne pas bloquer Redis.
|
|
func scanRedisKeys(pattern string) ([]string, error) {
|
|
var all []string
|
|
cursor := uint64(0)
|
|
for {
|
|
batch, next, err := Redis.Scan(RedisCtx, cursor, pattern, 200).Result()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
all = append(all, batch...)
|
|
cursor = next
|
|
if cursor == 0 {
|
|
break
|
|
}
|
|
}
|
|
return all, nil
|
|
}
|