From af22430101f7950d978e223e37cafef139adba99 Mon Sep 17 00:00:00 2001 From: shanshanzhong Date: Sun, 24 May 2026 23:42:37 -0700 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E8=AE=A2=E5=8D=95=20c?= =?UTF-8?q?laim=20=E6=9C=BA=E5=88=B6=E7=8A=B6=E6=80=81=E5=86=B2=E7=AA=81?= =?UTF-8?q?=E4=B8=8E=20stuck=20=E8=AE=A2=E5=8D=95=E6=81=A2=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 将 OrderStatusClaimed 从 4 改为 6,与 OrderStatusFailed(4) 区分 - releaseClaim 改为返回 error,失败时不再静默吞掉 - claimAndGetOrder 对 status=claimed 返回可重试错误而非静默跳过 - ProcessTask 区分 claimed stuck 错误(触发重试)和其他非 paid 状态(跳过) - finalizeCouponAndOrder 使用直接 DB 更新 claimed→finished,绕过 UpdateOrderStatus 的 status 5 无法通过该条件) 并显式删除 Redis 缓存避免缓存脏读 - 新增 StuckOrderRecoveryLogic:每 10 分钟扫描超时 claimed 订单, 重置为 paid 并重新入队激活任务 Co-authored-by: multica-agent --- queue/handler/routes.go | 3 + queue/logic/order/activateOrderLogic.go | 58 +++++++++++-- queue/logic/order/stuckOrderRecoveryLogic.go | 86 ++++++++++++++++++++ queue/types/scheduler.go | 1 + scheduler/scheduler.go | 6 ++ 5 files changed, 146 insertions(+), 8 deletions(-) create mode 100644 queue/logic/order/stuckOrderRecoveryLogic.go diff --git a/queue/handler/routes.go b/queue/handler/routes.go index ab252a7..9212a2d 100644 --- a/queue/handler/routes.go +++ b/queue/handler/routes.go @@ -48,4 +48,7 @@ func RegisterHandlers(mux *asynq.ServeMux, serverCtx *svc.ServiceContext) { // Apple IAP 对账(第二层:5min 扫描 + 第三层:日终全量) mux.Handle(types.SchedulerIAPReconcile, iapLogic.NewReconcileLogic(serverCtx)) mux.Handle(types.SchedulerIAPDailyReconcile, iapLogic.NewDailyReconcileLogic(serverCtx)) + + // Stuck order recovery: reset claimed orders that timed out back to paid + mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx)) } diff --git a/queue/logic/order/activateOrderLogic.go b/queue/logic/order/activateOrderLogic.go index ef0e24c..093d547 100644 --- a/queue/logic/order/activateOrderLogic.go +++ b/queue/logic/order/activateOrderLogic.go @@ -7,6 +7,7 @@ import ( "encoding/json" "errors" "fmt" + "strings" "time" "github.com/perfect-panel/server/internal/logic/admin/group" @@ -46,7 +47,7 @@ const ( OrderStatusPaid = 2 // Order paid and ready for processing OrderStatusClose = 3 // Order closed/cancelled OrderStatusFailed = 4 // Order processing failed - OrderStatusClaimed = 4 // Internal transient claim while a worker processes the order + OrderStatusClaimed = 6 // Internal transient claim while a worker processes the order OrderStatusFinished = 5 // Order successfully completed ) @@ -94,6 +95,11 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task) if err != nil { // 如果订单不存在或状态不对,不重试 if errors.Is(err, ErrInvalidOrderStatus) { + if strings.Contains(err.Error(), "stuck in claimed") { + logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单卡在 claimed,将重试", + logger.Field("order_no", payload.OrderNo)) + return err + } logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过", logger.Field("order_no", payload.OrderNo)) return nil @@ -116,7 +122,12 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task) ) if err = l.processOrderByType(ctx, orderInfo, payload.IAPExpireAt); err != nil { - l.releaseClaim(ctx, orderInfo.OrderNo) + if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil { + logger.WithContext(ctx).Error("[ActivateOrderLogic] releaseClaim also failed, stuck recovery will handle", + logger.Field("order_no", orderInfo.OrderNo), + logger.Field("release_error", releaseErr.Error()), + ) + } logger.WithContext(ctx).Error("[ActivateOrderLogic] 处理订单失败,将重试", logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_type", orderInfo.Type), @@ -125,7 +136,12 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task) } if err = l.reconcilePostOrderSubscriptions(ctx, orderInfo); err != nil { - l.releaseClaim(ctx, orderInfo.OrderNo) + if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil { + logger.WithContext(ctx).Error("[ActivateOrderLogic] releaseClaim also failed, stuck recovery will handle", + logger.Field("order_no", orderInfo.OrderNo), + logger.Field("release_error", releaseErr.Error()), + ) + } logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单订阅兜底合并失败,将重试", logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_type", orderInfo.Type), @@ -176,6 +192,14 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin return nil, nil } + if orderInfo.Status == OrderStatusClaimed { + logger.WithContext(ctx).Error("Order stuck in claimed status", + logger.Field("order_no", orderInfo.OrderNo), + logger.Field("status", orderInfo.Status), + ) + return nil, fmt.Errorf("order %s stuck in claimed status: %w", orderNo, ErrInvalidOrderStatus) + } + if orderInfo.Status != OrderStatusPaid { logger.WithContext(ctx).Error("Order status error", logger.Field("order_no", orderInfo.OrderNo), @@ -201,7 +225,7 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin return &orderInfo, nil } -func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) { +func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) error { if err := l.svc.DB.WithContext(ctx). Model(&order.Order{}). Where("order_no = ? AND status = ?", orderNo, OrderStatusClaimed). @@ -210,7 +234,9 @@ func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) { logger.Field("error", err.Error()), logger.Field("order_no", orderNo), ) + return fmt.Errorf("release claim failed for order %s: %w", orderNo, err) } + return nil } // processOrderByType routes order processing based on the order type @@ -563,10 +589,26 @@ func (l *ActivateOrderLogic) finalizeCouponAndOrder(ctx context.Context, orderIn } } - // Update order status using state-guarded UpdateOrderStatus to prevent double finalization - if err := l.svc.OrderModel.UpdateOrderStatus(ctx, orderInfo.OrderNo, OrderStatusFinished); err != nil { - logger.WithContext(ctx).Error("Update order status failed", - logger.Field("error", err.Error()), + // Direct update from claimed(6) → finished(5), bypassing the model's status 5 would never match. + result := l.svc.DB.WithContext(ctx). + Model(&order.Order{}). + Where("order_no = ? AND status = ?", orderInfo.OrderNo, OrderStatusClaimed). + Update("status", OrderStatusFinished) + if result.Error != nil { + logger.WithContext(ctx).Error("Update order status from claimed to finished failed", + logger.Field("error", result.Error.Error()), + logger.Field("order_no", orderInfo.OrderNo), + ) + } + // Invalidate cache entries; key format matches order/default.go cacheOrderIdPrefix / cacheOrderNoPrefix. + cacheKeys := []string{ + fmt.Sprintf("cache:order:id:%d", orderInfo.Id), + fmt.Sprintf("cache:order:no:%s", orderInfo.OrderNo), + } + if delErr := l.svc.Redis.Del(ctx, cacheKeys...).Err(); delErr != nil { + logger.WithContext(ctx).Error("Failed to invalidate order cache after status update", + logger.Field("error", delErr.Error()), logger.Field("order_no", orderInfo.OrderNo), ) } diff --git a/queue/logic/order/stuckOrderRecoveryLogic.go b/queue/logic/order/stuckOrderRecoveryLogic.go new file mode 100644 index 0000000..7cd9ee2 --- /dev/null +++ b/queue/logic/order/stuckOrderRecoveryLogic.go @@ -0,0 +1,86 @@ +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" +) + +const stuckClaimTimeout = 10 * time.Minute + +// StuckOrderRecoveryLogic scans orders stuck in claimed(6) status and resets them to paid(2) +// so that asynq retry can re-claim and process them. +type StuckOrderRecoveryLogic struct { + svc *svc.ServiceContext +} + +func NewStuckOrderRecoveryLogic(svc *svc.ServiceContext) *StuckOrderRecoveryLogic { + return &StuckOrderRecoveryLogic{svc: svc} +} + +func (l *StuckOrderRecoveryLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error { + cutoff := time.Now().Add(-stuckClaimTimeout) + + 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, resetting to paid", + logger.Field("count", len(stuckOrders)), + logger.Field("order_nos", orderNos), + ) + + result := l.svc.DB.WithContext(ctx). + Model(&order.Order{}). + Where("status = ? AND updated_at < ?", OrderStatusClaimed, cutoff). + Update("status", OrderStatusPaid) + if result.Error != nil { + logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to reset stuck orders", + logger.Field("error", result.Error.Error()), + ) + return result.Error + } + + for i := range stuckOrders { + ord := &stuckOrders[i] + payload, err := json.Marshal(queueTypes.ForthwithActivateOrderPayload{OrderNo: ord.OrderNo}) + if err != nil { + logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to marshal payload", + logger.Field("order_no", ord.OrderNo), + logger.Field("error", err.Error()), + ) + continue + } + task := asynq.NewTask(queueTypes.ForthwithActivateOrder, payload, asynq.MaxRetry(5)) + if _, err := l.svc.Queue.EnqueueContext(ctx, task); err != nil { + logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to re-enqueue order", + logger.Field("order_no", ord.OrderNo), + logger.Field("error", err.Error()), + ) + } + } + + return nil +} diff --git a/queue/types/scheduler.go b/queue/types/scheduler.go index 8bf2631..7d2c431 100644 --- a/queue/types/scheduler.go +++ b/queue/types/scheduler.go @@ -7,4 +7,5 @@ const ( SchedulerTrafficStat = "scheduler:traffic:stat" SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单 SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账 + SchedulerStuckOrderRecovery = "scheduler:stuck:order:recovery" // 定时恢复卡在 claimed 状态的订单 ) diff --git a/scheduler/scheduler.go b/scheduler/scheduler.go index f554cc3..f2b16df 100644 --- a/scheduler/scheduler.go +++ b/scheduler/scheduler.go @@ -64,6 +64,12 @@ func (m *Service) Start() { logger.Errorf("register iap daily reconcile task failed: %s", err.Error()) } + // schedule stuck order recovery: every 10 minutes + stuckOrderTask := asynq.NewTask(types.SchedulerStuckOrderRecovery, nil) + if _, err := m.server.Register("@every 10m", stuckOrderTask, asynq.MaxRetry(1)); err != nil { + logger.Errorf("register stuck order recovery task failed: %s", err.Error()) + } + if err := m.server.Run(); err != nil { logger.Errorf("run scheduler failed: %s", err.Error()) }