fix: 修复订单 claim 机制状态冲突与 stuck 订单恢复
Build docker and publish / build (20.15.1) (pull_request) Failing after 8m15s
Build docker and publish / build (20.15.1) (pull_request) Failing after 8m15s
- 将 OrderStatusClaimed 从 4 改为 6,与 OrderStatusFailed(4) 区分 - releaseClaim 改为返回 error,失败时不再静默吞掉 - claimAndGetOrder 对 status=claimed 返回可重试错误而非静默跳过 - ProcessTask 区分 claimed stuck 错误(触发重试)和其他非 paid 状态(跳过) - finalizeCouponAndOrder 使用直接 DB 更新 claimed→finished,绕过 UpdateOrderStatus 的 status<target 守卫(6 > 5 无法通过该条件) 并显式删除 Redis 缓存避免缓存脏读 - 新增 StuckOrderRecoveryLogic:每 10 分钟扫描超时 claimed 订单, 重置为 paid 并重新入队激活任务 Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
@@ -48,4 +48,7 @@ func RegisterHandlers(mux *asynq.ServeMux, serverCtx *svc.ServiceContext) {
|
|||||||
// Apple IAP 对账(第二层:5min 扫描 + 第三层:日终全量)
|
// Apple IAP 对账(第二层:5min 扫描 + 第三层:日终全量)
|
||||||
mux.Handle(types.SchedulerIAPReconcile, iapLogic.NewReconcileLogic(serverCtx))
|
mux.Handle(types.SchedulerIAPReconcile, iapLogic.NewReconcileLogic(serverCtx))
|
||||||
mux.Handle(types.SchedulerIAPDailyReconcile, iapLogic.NewDailyReconcileLogic(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))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/perfect-panel/server/internal/logic/admin/group"
|
"github.com/perfect-panel/server/internal/logic/admin/group"
|
||||||
@@ -46,7 +47,7 @@ const (
|
|||||||
OrderStatusPaid = 2 // Order paid and ready for processing
|
OrderStatusPaid = 2 // Order paid and ready for processing
|
||||||
OrderStatusClose = 3 // Order closed/cancelled
|
OrderStatusClose = 3 // Order closed/cancelled
|
||||||
OrderStatusFailed = 4 // Order processing failed
|
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
|
OrderStatusFinished = 5 // Order successfully completed
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -94,6 +95,11 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
// 如果订单不存在或状态不对,不重试
|
// 如果订单不存在或状态不对,不重试
|
||||||
if errors.Is(err, ErrInvalidOrderStatus) {
|
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.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过",
|
||||||
logger.Field("order_no", payload.OrderNo))
|
logger.Field("order_no", payload.OrderNo))
|
||||||
return nil
|
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 {
|
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.WithContext(ctx).Error("[ActivateOrderLogic] 处理订单失败,将重试",
|
||||||
logger.Field("order_no", orderInfo.OrderNo),
|
logger.Field("order_no", orderInfo.OrderNo),
|
||||||
logger.Field("order_type", orderInfo.Type),
|
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 {
|
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.WithContext(ctx).Error("[ActivateOrderLogic] 订单订阅兜底合并失败,将重试",
|
||||||
logger.Field("order_no", orderInfo.OrderNo),
|
logger.Field("order_no", orderInfo.OrderNo),
|
||||||
logger.Field("order_type", orderInfo.Type),
|
logger.Field("order_type", orderInfo.Type),
|
||||||
@@ -176,6 +192,14 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
|
|||||||
return nil, nil
|
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 {
|
if orderInfo.Status != OrderStatusPaid {
|
||||||
logger.WithContext(ctx).Error("Order status error",
|
logger.WithContext(ctx).Error("Order status error",
|
||||||
logger.Field("order_no", orderInfo.OrderNo),
|
logger.Field("order_no", orderInfo.OrderNo),
|
||||||
@@ -201,7 +225,7 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
|
|||||||
return &orderInfo, nil
|
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).
|
if err := l.svc.DB.WithContext(ctx).
|
||||||
Model(&order.Order{}).
|
Model(&order.Order{}).
|
||||||
Where("order_no = ? AND status = ?", orderNo, OrderStatusClaimed).
|
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("error", err.Error()),
|
||||||
logger.Field("order_no", orderNo),
|
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
|
// 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
|
// Direct update from claimed(6) → finished(5), bypassing the model's status<target guard.
|
||||||
if err := l.svc.OrderModel.UpdateOrderStatus(ctx, orderInfo.OrderNo, OrderStatusFinished); err != nil {
|
// UpdateOrderStatus uses WHERE status < ?, so status=6 > 5 would never match.
|
||||||
logger.WithContext(ctx).Error("Update order status failed",
|
result := l.svc.DB.WithContext(ctx).
|
||||||
logger.Field("error", err.Error()),
|
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),
|
logger.Field("order_no", orderInfo.OrderNo),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -7,4 +7,5 @@ const (
|
|||||||
SchedulerTrafficStat = "scheduler:traffic:stat"
|
SchedulerTrafficStat = "scheduler:traffic:stat"
|
||||||
SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单
|
SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单
|
||||||
SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账
|
SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账
|
||||||
|
SchedulerStuckOrderRecovery = "scheduler:stuck:order:recovery" // 定时恢复卡在 claimed 状态的订单
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -64,6 +64,12 @@ func (m *Service) Start() {
|
|||||||
logger.Errorf("register iap daily reconcile task failed: %s", err.Error())
|
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 {
|
if err := m.server.Run(); err != nil {
|
||||||
logger.Errorf("run scheduler failed: %s", err.Error())
|
logger.Errorf("run scheduler failed: %s", err.Error())
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user