From 294bdf457724dcc8cf6d5309097d1364147a7422 Mon Sep 17 00:00:00 2001 From: shanshanzhong Date: Sun, 24 May 2026 23:42:41 -0700 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E8=AE=A2=E5=8D=95?= =?UTF-8?q?=E7=8A=B6=E6=80=81=E6=9C=BA=20claim=20=E6=9C=BA=E5=88=B6?= =?UTF-8?q?=E7=9A=84=E4=B8=89=E4=B8=AA=20Bug=20=E5=B9=B6=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=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 的值冲突 - finalizeCouponAndOrder 改用直接 DB 更新(WHERE status=6→SET status=5), 绕过 UpdateOrderStatus 的 status Co-authored-by: multica-agent --- queue/handler/routes.go | 3 + queue/logic/order/activateOrderLogic.go | 57 ++++++++-- queue/logic/order/stuckOrderRecoveryLogic.go | 104 +++++++++++++++++++ queue/types/scheduler.go | 5 +- scheduler/scheduler.go | 6 ++ 5 files changed, 164 insertions(+), 11 deletions(-) create mode 100644 queue/logic/order/stuckOrderRecoveryLogic.go diff --git a/queue/handler/routes.go b/queue/handler/routes.go index ab252a7..03ae41d 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 + mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx)) } diff --git a/queue/logic/order/activateOrderLogic.go b/queue/logic/order/activateOrderLogic.go index ef0e24c..84ec550 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 ) @@ -92,8 +93,12 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task) orderInfo, err := l.claimAndGetOrder(ctx, payload.OrderNo) 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 // 返回错误触发 asynq 重试 + } logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过", logger.Field("order_no", payload.OrderNo)) return nil @@ -116,7 +121,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 +135,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 +191,15 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin return nil, nil } + // Detect stuck claimed order — return retryable error so asynq re-tries + 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,14 +589,27 @@ 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", + // UpdateOrderStatus uses WHERE status < target, which blocks claimed(6)→finished(5). + // Use a direct update matching the exact claimed status, then update the full record + // via the model layer to properly invalidate the cache. + 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 order cache regardless of whether the DB update succeeded + orderInfo.Status = OrderStatusFinished + if err := l.svc.OrderModel.Update(ctx, orderInfo); err != nil { + logger.WithContext(ctx).Error("Update order cache after finalization failed", logger.Field("error", err.Error()), logger.Field("order_no", orderInfo.OrderNo), ) } - orderInfo.Status = OrderStatusFinished commonLogic.SubscriptionTraceInfo(logger.WithContext(ctx), commonLogic.SubscriptionTraceFlowOrder, "order_status_finished", "[SubscriptionFlow] order status updated to finished", commonLogic.OrderTraceFields(orderInfo)..., diff --git a/queue/logic/order/stuckOrderRecoveryLogic.go b/queue/logic/order/stuckOrderRecoveryLogic.go new file mode 100644 index 0000000..20e5d40 --- /dev/null +++ b/queue/logic/order/stuckOrderRecoveryLogic.go @@ -0,0 +1,104 @@ +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 +} diff --git a/queue/types/scheduler.go b/queue/types/scheduler.go index 8bf2631..5e665fa 100644 --- a/queue/types/scheduler.go +++ b/queue/types/scheduler.go @@ -5,6 +5,7 @@ const ( SchedulerTotalServerData = "scheduler:total:server" SchedulerResetTraffic = "scheduler:reset:traffic" SchedulerTrafficStat = "scheduler:traffic:stat" - SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单 - SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账 + 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()) }