package orderLogic import ( "context" "encoding/json" "time" "github.com/hibiken/asynq" "github.com/perfect-panel/server/internal/model/order" "github.com/perfect-panel/server/internal/svc" "github.com/perfect-panel/server/pkg/logger" queueTypes "github.com/perfect-panel/server/queue/types" ) // StuckOrderRecoveryLogic scans orders stuck in claimed status and re-queues them for processing. type StuckOrderRecoveryLogic struct { svc *svc.ServiceContext } func NewStuckOrderRecoveryLogic(svc *svc.ServiceContext) *StuckOrderRecoveryLogic { return &StuckOrderRecoveryLogic{svc: svc} } // ProcessTask scans for orders stuck in claimed status for over 10 minutes, // resets them to paid, and re-enqueues the activate task so they are retried // independently of asynq's original retry counter (which may be exhausted). func (l *StuckOrderRecoveryLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error { cutoff := time.Now().Add(-10 * time.Minute) var stuckOrders []order.Order if err := l.svc.DB.WithContext(ctx). Model(&order.Order{}). Where("status = ? AND updated_at < ?", OrderStatusClaimed, cutoff). Find(&stuckOrders).Error; err != nil { logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to query stuck orders", logger.Field("error", err.Error()), ) return err } if len(stuckOrders) == 0 { return nil } orderNos := make([]string, 0, len(stuckOrders)) for i := range stuckOrders { orderNos = append(orderNos, stuckOrders[i].OrderNo) } logger.WithContext(ctx).Error("[StuckOrderRecovery] Found stuck claimed orders, recovering", logger.Field("count", len(stuckOrders)), logger.Field("order_nos", orderNos), ) for i := range stuckOrders { o := &stuckOrders[i] result := l.svc.DB.WithContext(ctx). Model(&order.Order{}). Where("order_no = ? AND status = ?", o.OrderNo, OrderStatusClaimed). Update("status", OrderStatusPaid) if result.Error != nil { logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to reset order status", logger.Field("order_no", o.OrderNo), logger.Field("error", result.Error.Error()), ) continue } if result.RowsAffected == 0 { // Another process already handled this order continue } // Invalidate order cache o.Status = OrderStatusPaid if err := l.svc.OrderModel.Update(ctx, o); err != nil { logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to update order cache", logger.Field("order_no", o.OrderNo), logger.Field("error", err.Error()), ) } // Re-enqueue activate task so the order gets processed regardless of asynq retry state payload, err := json.Marshal(queueTypes.ForthwithActivateOrderPayload{OrderNo: o.OrderNo}) if err != nil { logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to marshal task payload", logger.Field("order_no", o.OrderNo), logger.Field("error", err.Error()), ) continue } if _, err = l.svc.Queue.EnqueueContext(ctx, asynq.NewTask(queueTypes.ForthwithActivateOrder, payload)); err != nil { logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to re-enqueue activate task", logger.Field("order_no", o.OrderNo), logger.Field("error", err.Error()), ) } else { logger.WithContext(ctx).Info("[StuckOrderRecovery] Re-enqueued activate task", logger.Field("order_no", o.OrderNo), ) } } return nil }