391 lines
11 KiB
Go
391 lines
11 KiB
Go
// ============================================
|
|
// db/cancel_commands_db.go
|
|
// FONCTIONS DB ATOMIQUES POUR L'ANNULATION
|
|
// VERSION 100% SÉCURISÉE - FIX ETA CHECK
|
|
// ============================================
|
|
|
|
package db
|
|
|
|
import (
|
|
"database/sql"
|
|
"fmt"
|
|
"log"
|
|
"slices"
|
|
"time"
|
|
)
|
|
|
|
// ============================================
|
|
// ANNULATION ATOMIQUE
|
|
// ============================================
|
|
|
|
func (d *Database) CancelCommandAtomic(commandID int, username, reason string, force bool) (int, error) {
|
|
log.Printf("🔒 [CancelAtomic] START - cmd=%d, user=%s, force=%v", commandID, username, force)
|
|
|
|
tx, err := d.Begin()
|
|
if err != nil {
|
|
return 0, fmt.Errorf("erreur transaction: %w", err)
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
var currentStatus, cmdUsername, livreurAssign string
|
|
err = tx.QueryRow(`
|
|
SELECT status, username, COALESCE(livreur_assign, '')
|
|
FROM commandes
|
|
WHERE id = $1
|
|
FOR UPDATE
|
|
`, commandID).Scan(¤tStatus, &cmdUsername, &livreurAssign)
|
|
|
|
if err == sql.ErrNoRows {
|
|
return 0, fmt.Errorf("commande non trouvée")
|
|
}
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
log.Printf("📋 [CancelAtomic] Trouvée - status=%s, owner=%s, livreur=%s", currentStatus, cmdUsername, livreurAssign)
|
|
|
|
if cmdUsername != username {
|
|
return 0, fmt.Errorf("commande ne vous appartient pas")
|
|
}
|
|
|
|
nonCancellableStatuses := []string{"livre", "approved", "cancelled", "disabled"}
|
|
if slices.Contains(nonCancellableStatuses, currentStatus) {
|
|
return 0, fmt.Errorf("impossible d'annuler")
|
|
}
|
|
|
|
isLateCancel := false
|
|
if livreurAssign != "" {
|
|
if currentStatus == "en_route" || currentStatus == "arrived" {
|
|
isLateCancel = true
|
|
log.Printf("⚠️ [CancelAtomic] Annulation TARDIVE détectée - Statut: %s", currentStatus)
|
|
} else if d.CheckCommandETAExistsAndValid(commandID) {
|
|
isLateCancel = true
|
|
log.Printf("⚠️ [CancelAtomic] Annulation TARDIVE détectée - ETA définie")
|
|
} else {
|
|
log.Printf("✅ [CancelAtomic] Annulation SANS PÉNALITÉ - Statut: %s, Pas d'ETA valide", currentStatus)
|
|
}
|
|
} else {
|
|
log.Printf("✅ [CancelAtomic] Annulation SANS PÉNALITÉ - Aucun livreur assigné")
|
|
}
|
|
|
|
if isLateCancel && !force {
|
|
return 0, fmt.Errorf("confirmation requise")
|
|
}
|
|
|
|
result, err := tx.Exec(`
|
|
UPDATE commandes
|
|
SET status = 'cancelled', updated_at = CURRENT_TIMESTAMP
|
|
WHERE id = $1 AND status = $2 AND username = $3
|
|
`, commandID, currentStatus, username)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
if rows, _ := result.RowsAffected(); rows == 0 {
|
|
return 0, fmt.Errorf("commande déjà modifiée")
|
|
}
|
|
log.Printf("✅ [CancelAtomic] Statut mis à jour: %s → cancelled", currentStatus)
|
|
|
|
if _, err = tx.Exec(`
|
|
UPDATE products p
|
|
SET stock = stock + ci.quantite, updated_at = CURRENT_TIMESTAMP
|
|
FROM command_items ci
|
|
WHERE ci.command_id = $1 AND ci.product_id = p.id
|
|
`, commandID); err != nil {
|
|
log.Printf("⚠️ [CancelAtomic] Erreur remboursement stock: %v", err)
|
|
} else {
|
|
log.Printf("✅ [CancelAtomic] Stock remboursé")
|
|
}
|
|
|
|
penalty := 0
|
|
if isLateCancel {
|
|
log.Printf("⚠️ [CancelAtomic] Annulation tardive confirmée - Application pénalité")
|
|
penalty, _ = d.CalculateCancellationPenalty(username)
|
|
if _, err = tx.Exec(`
|
|
UPDATE clients
|
|
SET amende = amende + $1,
|
|
cancellations_count = COALESCE(cancellations_count, 0) + 1,
|
|
updated_at = CURRENT_TIMESTAMP
|
|
WHERE username = $2
|
|
`, penalty, username); err != nil {
|
|
log.Printf("❌ [CancelAtomic] Erreur pénalité: %v", err)
|
|
} else {
|
|
log.Printf("⚠️ [CancelAtomic] Pénalité: %d appliquée à %s", penalty, username)
|
|
}
|
|
}
|
|
|
|
if _, err = tx.Exec(`
|
|
INSERT INTO command_logs (command_id, status, message, author, created_at)
|
|
VALUES ($1, 'cancelled', $2, $3, CURRENT_TIMESTAMP)
|
|
`, commandID, fmt.Sprintf("Annulée par %s - Raison: %s", username, reason), username); err != nil {
|
|
log.Printf("⚠️ [CancelAtomic] Erreur log: %v", err)
|
|
}
|
|
|
|
if err := tx.Commit(); err != nil {
|
|
return 0, fmt.Errorf("erreur commit: %w", err)
|
|
}
|
|
|
|
if livreurAssign != "" {
|
|
go func() {
|
|
if err := d.CleanupCompletedCommandFromQueue(commandID, livreurAssign); err != nil {
|
|
log.Printf("⚠️ [CancelAtomic] Erreur cleanup queue: %v", err)
|
|
}
|
|
}()
|
|
}
|
|
go func() {
|
|
Redis.Del(RedisCtx,
|
|
fmt.Sprintf("command:%d", commandID),
|
|
fmt.Sprintf("client:%s", username),
|
|
fmt.Sprintf("client:%s:commands", username),
|
|
)
|
|
}()
|
|
|
|
log.Printf("🎉 [CancelAtomic] SUCCÈS - Commande %d annulée", commandID)
|
|
return penalty, nil
|
|
}
|
|
|
|
// ============================================
|
|
// ✅ NOUVELLE FONCTION: CHECK ETA VALIDE
|
|
// ============================================
|
|
|
|
// CheckCommandETAExistsAndValid vérifie si une ETA RÉELLE existe (> 0 minutes, non expirée)
|
|
func (d *Database) CheckCommandETAExistsAndValid(commandID int) bool {
|
|
etaKey := fmt.Sprintf("command:eta:%d", commandID)
|
|
|
|
// Récupérer l'ETA depuis Redis
|
|
etaMinutesStr, err := Redis.Get(RedisCtx, etaKey).Result()
|
|
if err != nil {
|
|
log.Printf("⚠️ [CheckETA] Pas d'ETA trouvée pour cmd %d", commandID)
|
|
return false
|
|
}
|
|
|
|
// Parser l'ETA
|
|
var etaMinutes int
|
|
_, err = fmt.Sscanf(etaMinutesStr, "%d", &etaMinutes)
|
|
if err != nil || etaMinutes <= 0 {
|
|
log.Printf("⚠️ [CheckETA] ETA invalide pour cmd %d: %s", commandID, etaMinutesStr)
|
|
return false
|
|
}
|
|
|
|
// Vérifier le TTL (si l'ETA existe, elle doit avoir un TTL)
|
|
ttl, err := Redis.TTL(RedisCtx, etaKey).Result()
|
|
if err != nil || ttl <= 0 {
|
|
log.Printf("⚠️ [CheckETA] ETA expirée pour cmd %d", commandID)
|
|
return false
|
|
}
|
|
|
|
log.Printf("✅ [CheckETA] ETA valide trouvée pour cmd %d: %d min (TTL: %v)", commandID, etaMinutes, ttl)
|
|
return true
|
|
}
|
|
|
|
// ============================================
|
|
// SUPPRESSION ATOMIQUE
|
|
// ============================================
|
|
|
|
func (d *Database) DeleteCommandAtomic(commandID int, deletedBy, role string) error {
|
|
log.Printf("🔒 [DeleteAtomic] START - cmd=%d, by=%s (%s)", commandID, deletedBy, role)
|
|
|
|
// ✅ TRANSACTION
|
|
tx, err := d.Begin()
|
|
if err != nil {
|
|
return fmt.Errorf("erreur transaction: %w", err)
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
// ✅ SELECT FOR UPDATE
|
|
var currentStatus, cmdUsername, livreurAssign string
|
|
err = tx.QueryRow(`
|
|
SELECT status, username, COALESCE(livreur_assign, '')
|
|
FROM commandes
|
|
WHERE id = $1
|
|
FOR UPDATE
|
|
`, commandID).Scan(¤tStatus, &cmdUsername, &livreurAssign)
|
|
|
|
if err == sql.ErrNoRows {
|
|
return fmt.Errorf("commande non trouvée")
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
log.Printf("📋 [DeleteAtomic] Trouvée - status=%s, client=%s", currentStatus, cmdUsername)
|
|
|
|
// ✅ REMBOURSER STOCK ATOMIQUEMENT
|
|
_, err = tx.Exec(`
|
|
UPDATE products p
|
|
SET stock = stock + ci.quantite, updated_at = CURRENT_TIMESTAMP
|
|
FROM command_items ci
|
|
WHERE ci.command_id = $1 AND ci.product_id = p.id
|
|
`, commandID)
|
|
|
|
if err != nil {
|
|
log.Printf("⚠️ [DeleteAtomic] Erreur remboursement: %v", err)
|
|
} else {
|
|
log.Printf("✅ [DeleteAtomic] Stock remboursé")
|
|
}
|
|
|
|
// ✅ LOG AVANT SUPPRESSION
|
|
_, err = tx.Exec(`
|
|
INSERT INTO command_logs (command_id, status, message, author, created_at)
|
|
VALUES ($1, $2, $3, $4, CURRENT_TIMESTAMP)
|
|
`, commandID, "deleted",
|
|
fmt.Sprintf("Supprimée par %s (%s) - Ancien statut: %s", deletedBy, role, currentStatus),
|
|
deletedBy)
|
|
|
|
// ✅ SUPPRIMER ITEMS
|
|
_, err = tx.Exec(`DELETE FROM command_items WHERE command_id = $1`, commandID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// ✅ SUPPRIMER COMMANDE
|
|
result, err := tx.Exec(`DELETE FROM commandes WHERE id = $1`, commandID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
rows, _ := result.RowsAffected()
|
|
if rows == 0 {
|
|
return fmt.Errorf("commande non trouvée")
|
|
}
|
|
|
|
log.Printf("✅ [DeleteAtomic] Supprimée de la DB")
|
|
|
|
// ✅ COMMIT
|
|
if err := tx.Commit(); err != nil {
|
|
return fmt.Errorf("erreur commit: %w", err)
|
|
}
|
|
|
|
// ✅ NETTOYER QUEUES (async)
|
|
if livreurAssign != "" {
|
|
go d.RemoveCommandFromAllQueues(commandID, livreurAssign)
|
|
}
|
|
|
|
// ✅ INVALIDER CACHES (async)
|
|
go func() {
|
|
Redis.Del(RedisCtx,
|
|
fmt.Sprintf("command:%d", commandID),
|
|
fmt.Sprintf("client:%s", cmdUsername),
|
|
fmt.Sprintf("client:%s:commands", cmdUsername),
|
|
)
|
|
}()
|
|
|
|
log.Printf("🎉 [DeleteAtomic] SUCCÈS - Commande %d supprimée", commandID)
|
|
return nil
|
|
}
|
|
|
|
// ============================================
|
|
// FONCTIONS HELPERS (déjà sécurisées)
|
|
// ============================================
|
|
|
|
func (d *Database) GetCommandPositionInQueue(livreurUsername string, commandID int) (int, error) {
|
|
queueKey := fmt.Sprintf("queue:deliveryman:%s", livreurUsername)
|
|
commandIDStr := fmt.Sprintf("%d", commandID)
|
|
|
|
rank, err := Redis.ZRank(RedisCtx, queueKey, commandIDStr).Result()
|
|
if err != nil {
|
|
return 0, fmt.Errorf("commande non trouvée dans la queue")
|
|
}
|
|
|
|
return int(rank) + 1, nil
|
|
}
|
|
|
|
func (d *Database) GetCancelledCommands(username string, limit int) ([]map[string]interface{}, error) {
|
|
query := `
|
|
SELECT id, username, status, adresse, total_prix, created_at, updated_at
|
|
FROM commandes
|
|
WHERE status = 'cancelled'
|
|
`
|
|
|
|
args := []interface{}{}
|
|
argPos := 1
|
|
|
|
if username != "" {
|
|
query += fmt.Sprintf(" AND username = $%d", argPos)
|
|
args = append(args, username)
|
|
argPos++
|
|
}
|
|
|
|
query += " ORDER BY updated_at DESC"
|
|
|
|
if limit > 0 {
|
|
query += fmt.Sprintf(" LIMIT $%d", argPos)
|
|
args = append(args, limit)
|
|
}
|
|
|
|
rows, err := d.Query(query, args...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("erreur récupération: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
var commands []map[string]interface{}
|
|
for rows.Next() {
|
|
var id int
|
|
var username, status, adresse string
|
|
var totalPrix float64
|
|
var createdAt, updatedAt time.Time
|
|
|
|
err := rows.Scan(&id, &username, &status, &adresse, &totalPrix, &createdAt, &updatedAt)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
commands = append(commands, map[string]interface{}{
|
|
"id": id,
|
|
"username": username,
|
|
"status": status,
|
|
"adresse": adresse,
|
|
"total_prix": totalPrix,
|
|
"created_at": createdAt,
|
|
"updated_at": updatedAt,
|
|
})
|
|
}
|
|
|
|
return commands, nil
|
|
}
|
|
|
|
// ============================================
|
|
// GESTION DES PÉNALITÉS
|
|
// ============================================
|
|
|
|
// AddClientPenalty ajoute une pénalité à un client
|
|
func (d *Database) AddClientPenalty(username string, points int) error {
|
|
log.Printf("⚠️ [AddPenalty] Ajout pénalité: %d points pour client %s", points, username)
|
|
|
|
// ✅ Validation
|
|
if points <= 0 {
|
|
return fmt.Errorf("points invalides: %d", points)
|
|
}
|
|
|
|
if username == "" {
|
|
return fmt.Errorf("username vide")
|
|
}
|
|
|
|
// ✅ UPDATE dans PostgreSQL
|
|
query := `UPDATE clients
|
|
SET amende = amende + $1, updated_at = CURRENT_TIMESTAMP
|
|
WHERE username = $2`
|
|
|
|
result, err := d.Exec(query, points, username)
|
|
if err != nil {
|
|
log.Printf("❌ [AddPenalty] Erreur UPDATE: %v", err)
|
|
return fmt.Errorf("erreur ajout pénalité: %w", err)
|
|
}
|
|
|
|
rowsAffected, err := result.RowsAffected()
|
|
if err != nil {
|
|
return fmt.Errorf("erreur vérification: %w", err)
|
|
}
|
|
if rowsAffected == 0 {
|
|
return fmt.Errorf("client non trouvé: %s", username)
|
|
}
|
|
|
|
log.Printf("✅ [AddPenalty] Pénalité ajoutée: +%d points pour %s", points, username)
|
|
|
|
// ✅ Invalider le cache Redis du client
|
|
cacheKey := fmt.Sprintf("client:%s", username)
|
|
Redis.Del(RedisCtx, cacheKey)
|
|
|
|
return nil
|
|
}
|