Compare commits

..

1 Commits

Author SHA1 Message Date
shanshanzhong147 28ad97def4 fix: persist guest/redemption activation context to DB to survive Redis TTL expiry (Bug 1 + Bug 10)
- Add `activation_context` TEXT column to `order` table (migration 02150)
- purchaseLogic: write TemporaryOrderInfo JSON to order.ActivationContext in the same
  insert transaction; Redis write is now best-effort (non-fatal)
- redeemCodeLogic: write redemption {type, redemption_code_id, unit_time, quantity} JSON
  to order.ActivationContext at order creation; Redis write is now best-effort (non-fatal)
- getTempOrderInfo: on Redis miss, fall back to order.ActivationContext from DB;
  logs CRITICAL and returns error if both are missing (old orders with no DB record)
- RedemptionActivate: on Redis miss, fall back to order.ActivationContext from DB;
  same CRITICAL log path for legacy orders

Co-authored-by: multica-agent <github@multica.ai>
2026-05-25 00:13:59 -07:00
10 changed files with 141 additions and 225 deletions
@@ -0,0 +1,2 @@
-- Remove activation_context column from order table
ALTER TABLE `order` DROP COLUMN IF EXISTS `activation_context`;
@@ -0,0 +1,6 @@
-- Add activation_context column to order table for Redis fallback persistence (idempotent)
SET @col_exists = (SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'order' AND COLUMN_NAME = 'activation_context');
SET @sql = IF(@col_exists = 0, 'ALTER TABLE `order` ADD COLUMN `activation_context` TEXT DEFAULT NULL COMMENT ''Activation context JSON (guest/redemption info for DB fallback)'' AFTER `app_account_token`', 'SELECT 1');
PREPARE stmt FROM @sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
@@ -154,9 +154,12 @@ func (l *PurchaseLogic) Purchase(req *types.PortalPurchaseRequest) (resp *types.
} }
content, _ := tempOrder.Marshal() content, _ := tempOrder.Marshal()
if _, err = l.svcCtx.Redis.Set(l.ctx, fmt.Sprintf(constant.TempOrderCacheKey, orderInfo.OrderNo), string(content), CloseOrderTimeMinutes*time.Minute).Result(); err != nil { // Persist activation context to DB so the worker can recover if Redis TTL expires.
l.Errorw("[Purchase] Redis set error", logger.Field("error", err.Error()), logger.Field("order_no", orderInfo.OrderNo)) orderInfo.ActivationContext = string(content)
return err
// Write to Redis as a hot cache (best-effort; non-fatal on failure).
if _, redisErr := l.svcCtx.Redis.Set(l.ctx, fmt.Sprintf(constant.TempOrderCacheKey, orderInfo.OrderNo), string(content), CloseOrderTimeMinutes*time.Minute).Result(); redisErr != nil {
l.Infow("[Purchase] Redis set error (non-fatal, DB fallback available)", logger.Field("error", redisErr.Error()), logger.Field("order_no", orderInfo.OrderNo))
} }
l.Infow("[Purchase] Guest order", logger.Field("order_no", orderInfo.OrderNo), logger.Field("identifier", req.Identifier)) l.Infow("[Purchase] Guest order", logger.Field("order_no", orderInfo.OrderNo), logger.Field("identifier", req.Identifier))
@@ -169,7 +172,7 @@ func (l *PurchaseLogic) Purchase(req *types.PortalPurchaseRequest) (resp *types.
} }
} }
// save guest order // save guest order (activation_context is included)
if err = l.svcCtx.OrderModel.Insert(l.ctx, orderInfo, tx); err != nil { if err = l.svcCtx.OrderModel.Insert(l.ctx, orderInfo, tx); err != nil {
return err return err
} }
@@ -151,34 +151,48 @@ func (l *RedeemCodeLogic) RedeemCode(req *types.RedeemCodeRequest) (resp *types.
} }
// 创建Order记录 // 创建Order记录
redemptionContext := struct {
Type string `json:"type"`
RedemptionCodeId int64 `json:"redemption_code_id"`
UnitTime string `json:"unit_time"`
Quantity int64 `json:"quantity"`
}{
Type: "redemption",
RedemptionCodeId: redemptionCode.Id,
UnitTime: redemptionCode.UnitTime,
Quantity: redemptionCode.Quantity,
}
activationContextJSON, _ := json.Marshal(redemptionContext)
orderInfo := &order.Order{ orderInfo := &order.Order{
UserId: u.Id, UserId: u.Id,
OrderNo: tool.GenerateTradeNo(), OrderNo: tool.GenerateTradeNo(),
Type: 5, // 兑换类型 Type: 5, // 兑换类型
Quantity: redemptionCode.Quantity, Quantity: redemptionCode.Quantity,
Price: 0, // 兑换无价格 Price: 0, // 兑换无价格
Amount: 0, // 兑换无金额 Amount: 0, // 兑换无金额
Discount: 0, Discount: 0,
GiftAmount: 0, GiftAmount: 0,
Coupon: "", Coupon: "",
CouponDiscount: 0, CouponDiscount: 0,
PaymentId: 0, PaymentId: 0,
Method: "redemption", Method: "redemption",
FeeAmount: 0, FeeAmount: 0,
Commission: 0, Commission: 0,
Status: 2, // 直接设置为已支付 Status: 2, // 直接设置为已支付
SubscribeId: redemptionCode.SubscribePlan, SubscribeId: redemptionCode.SubscribePlan,
IsNew: isNew, IsNew: isNew,
ActivationContext: string(activationContextJSON),
} }
// 保存Order到数据库 // 保存Order到数据库activation_context 同步写入,作为 Redis 的持久化兜底)
err = l.svcCtx.OrderModel.Insert(l.ctx, orderInfo) err = l.svcCtx.OrderModel.Insert(l.ctx, orderInfo)
if err != nil { if err != nil {
l.Errorw("[RedeemCode] Create order failed", logger.Field("error", err.Error())) l.Errorw("[RedeemCode] Create order failed", logger.Field("error", err.Error()))
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "create order failed") return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "create order failed")
} }
// 缓存兑换码信息到Redis(供队列任务使用 // 缓存兑换码信息到Redis热缓存,供队列任务快速读取,非关键路径
cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo) cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo)
cacheData := map[string]interface{}{ cacheData := map[string]interface{}{
"redemption_code_id": redemptionCode.Id, "redemption_code_id": redemptionCode.Id,
@@ -186,16 +200,8 @@ func (l *RedeemCodeLogic) RedeemCode(req *types.RedeemCodeRequest) (resp *types.
"quantity": redemptionCode.Quantity, "quantity": redemptionCode.Quantity,
} }
jsonData, _ := json.Marshal(cacheData) jsonData, _ := json.Marshal(cacheData)
err = l.svcCtx.Redis.Set(l.ctx, cacheKey, jsonData, 2*time.Hour).Err() if redisErr := l.svcCtx.Redis.Set(l.ctx, cacheKey, jsonData, 2*time.Hour).Err(); redisErr != nil {
if err != nil { l.Infow("[RedeemCode] Cache redemption data failed (non-fatal, DB fallback available)", logger.Field("error", redisErr.Error()))
l.Errorw("[RedeemCode] Cache redemption data failed", logger.Field("error", err.Error()))
// 缓存失败,删除已创建的Order避免孤儿记录
if delErr := l.svcCtx.OrderModel.Delete(l.ctx, orderInfo.Id); delErr != nil {
l.Errorw("[RedeemCode] Delete order failed after cache error",
logger.Field("order_id", orderInfo.Id),
logger.Field("error", delErr.Error()))
}
return nil, errors.Wrapf(xerr.NewErrCode(xerr.ERROR), "cache redemption data failed")
} }
// 触发队列任务 // 触发队列任务
+3 -2
View File
@@ -24,8 +24,9 @@ type Order struct {
Status uint8 `gorm:"type:tinyint(1);not null;default:1;comment:Order Status: 1: Pending, 2: Paid, 3:Close, 4: Failed, 5:Finished;"` Status uint8 `gorm:"type:tinyint(1);not null;default:1;comment:Order Status: 1: Pending, 2: Paid, 3:Close, 4: Failed, 5:Finished;"`
SubscribeId int64 `gorm:"type:bigint;not null;default:0;comment:Subscribe Id"` SubscribeId int64 `gorm:"type:bigint;not null;default:0;comment:Subscribe Id"`
SubscribeToken string `gorm:"type:varchar(255);default:null;comment:Renewal Subscribe Token"` SubscribeToken string `gorm:"type:varchar(255);default:null;comment:Renewal Subscribe Token"`
AppAccountToken string `gorm:"type:varchar(36);default:null;comment:Apple IAP App Account Token (UUID)"` AppAccountToken string `gorm:"type:varchar(36);default:null;comment:Apple IAP App Account Token (UUID)"`
IsNew bool `gorm:"type:tinyint(1);not null;default:0;comment:Is New Order"` ActivationContext string `gorm:"type:text;default:null;comment:Activation context JSON (guest/redemption info for DB fallback)"`
IsNew bool `gorm:"type:tinyint(1);not null;default:0;comment:Is New Order"`
CreatedAt time.Time `gorm:"<-:create;comment:Create Time"` CreatedAt time.Time `gorm:"<-:create;comment:Create Time"`
UpdatedAt time.Time `gorm:"comment:Update Time"` UpdatedAt time.Time `gorm:"comment:Update Time"`
} }
-3
View File
@@ -48,7 +48,4 @@ 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
mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx))
} }
+86 -74
View File
@@ -7,7 +7,6 @@ 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"
@@ -47,7 +46,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 = 6 // Internal transient claim while a worker processes the order OrderStatusClaimed = 4 // Internal transient claim while a worker processes the order
OrderStatusFinished = 5 // Order successfully completed OrderStatusFinished = 5 // Order successfully completed
) )
@@ -93,12 +92,8 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
orderInfo, err := l.claimAndGetOrder(ctx, payload.OrderNo) orderInfo, err := l.claimAndGetOrder(ctx, payload.OrderNo)
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 // 返回错误触发 asynq 重试
}
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
@@ -121,12 +116,7 @@ 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 {
if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil { l.releaseClaim(ctx, orderInfo.OrderNo)
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),
@@ -135,12 +125,7 @@ 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 {
if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil { l.releaseClaim(ctx, orderInfo.OrderNo)
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),
@@ -191,15 +176,6 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
return nil, nil 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 { 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),
@@ -225,7 +201,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) error { func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) {
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).
@@ -234,9 +210,7 @@ func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) e
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
@@ -589,27 +563,14 @@ func (l *ActivateOrderLogic) finalizeCouponAndOrder(ctx context.Context, orderIn
} }
} }
// UpdateOrderStatus uses WHERE status < target, which blocks claimed(6)→finished(5). // Update order status using state-guarded UpdateOrderStatus to prevent double finalization
// Use a direct update matching the exact claimed status, then update the full record if err := l.svc.OrderModel.UpdateOrderStatus(ctx, orderInfo.OrderNo, OrderStatusFinished); err != nil {
// via the model layer to properly invalidate the cache. logger.WithContext(ctx).Error("Update order status failed",
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("error", err.Error()),
logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_no", orderInfo.OrderNo),
) )
} }
orderInfo.Status = OrderStatusFinished
commonLogic.SubscriptionTraceInfo(logger.WithContext(ctx), commonLogic.SubscriptionTraceFlowOrder, "order_status_finished", commonLogic.SubscriptionTraceInfo(logger.WithContext(ctx), commonLogic.SubscriptionTraceFlowOrder, "order_status_finished",
"[SubscriptionFlow] order status updated to finished", "[SubscriptionFlow] order status updated to finished",
commonLogic.OrderTraceFields(orderInfo)..., commonLogic.OrderTraceFields(orderInfo)...,
@@ -863,26 +824,52 @@ func (l *ActivateOrderLogic) createGuestUser(ctx context.Context, orderInfo *ord
return userInfo, nil return userInfo, nil
} }
// getTempOrderInfo retrieves temporary order information from Redis cache // getTempOrderInfo retrieves temporary order information from Redis cache with DB fallback.
func (l *ActivateOrderLogic) getTempOrderInfo(ctx context.Context, orderNo string) (*constant.TemporaryOrderInfo, error) { func (l *ActivateOrderLogic) getTempOrderInfo(ctx context.Context, orderNo string) (*constant.TemporaryOrderInfo, error) {
cacheKey := fmt.Sprintf(constant.TempOrderCacheKey, orderNo) cacheKey := fmt.Sprintf(constant.TempOrderCacheKey, orderNo)
data, err := l.svc.Redis.Get(ctx, cacheKey).Result() data, err := l.svc.Redis.Get(ctx, cacheKey).Result()
if err != nil { if err == nil {
logger.WithContext(ctx).Error("Get temp order cache failed", var tempOrder constant.TemporaryOrderInfo
logger.Field("error", err.Error()), if unmarshalErr := tempOrder.Unmarshal([]byte(data)); unmarshalErr != nil {
logger.Field("cache_key", cacheKey), logger.WithContext(ctx).Error("Unmarshal temp order cache failed",
logger.Field("error", unmarshalErr.Error()),
logger.Field("cache_key", cacheKey),
logger.Field("data", data),
)
return nil, unmarshalErr
}
return &tempOrder, nil
}
// Redis miss — fall back to DB activation_context field.
logger.WithContext(ctx).Infow("Redis cache miss for temp order, falling back to DB",
logger.Field("cache_key", cacheKey),
logger.Field("error", err.Error()),
)
orderInfo, dbErr := l.svc.OrderModel.FindOneByOrderNo(ctx, orderNo)
if dbErr != nil {
logger.WithContext(ctx).Error("DB fallback for temp order failed",
logger.Field("order_no", orderNo),
logger.Field("error", dbErr.Error()),
) )
return nil, err return nil, dbErr
}
if orderInfo.ActivationContext == "" {
logger.WithContext(ctx).Error("CRITICAL: activation_context missing in DB and Redis expired; cannot recover guest order",
logger.Field("order_no", orderNo),
)
return nil, errors.New("activation context not found: Redis expired and no DB fallback available")
} }
var tempOrder constant.TemporaryOrderInfo var tempOrder constant.TemporaryOrderInfo
if err = tempOrder.Unmarshal([]byte(data)); err != nil { if unmarshalErr := tempOrder.Unmarshal([]byte(orderInfo.ActivationContext)); unmarshalErr != nil {
logger.WithContext(ctx).Error("Unmarshal temp order cache failed", logger.WithContext(ctx).Error("Unmarshal DB activation_context failed",
logger.Field("error", err.Error()), logger.Field("order_no", orderNo),
logger.Field("cache_key", cacheKey), logger.Field("error", unmarshalErr.Error()),
logger.Field("data", data),
) )
return nil, err return nil, unmarshalErr
} }
return &tempOrder, nil return &tempOrder, nil
@@ -1782,25 +1769,50 @@ func (l *ActivateOrderLogic) RedemptionActivate(ctx context.Context, orderInfo *
return err return err
} }
// 3. 从Redis获取兑换码信息 // 3. 从Redis获取兑换码信息Redis缺失时回查DB
cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo)
data, err := l.svc.Redis.Get(ctx, cacheKey).Result()
if err != nil {
logger.WithContext(ctx).Error("Get redemption cache failed",
logger.Field("error", err.Error()),
logger.Field("cache_key", cacheKey),
)
return err
}
var redemptionData struct { var redemptionData struct {
RedemptionCodeId int64 `json:"redemption_code_id"` RedemptionCodeId int64 `json:"redemption_code_id"`
UnitTime string `json:"unit_time"` UnitTime string `json:"unit_time"`
Quantity int64 `json:"quantity"` Quantity int64 `json:"quantity"`
} }
if err = json.Unmarshal([]byte(data), &redemptionData); err != nil {
logger.WithContext(ctx).Error("Unmarshal redemption cache failed", logger.Field("error", err.Error())) cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo)
return err data, redisErr := l.svc.Redis.Get(ctx, cacheKey).Result()
if redisErr == nil {
if err = json.Unmarshal([]byte(data), &redemptionData); err != nil {
logger.WithContext(ctx).Error("Unmarshal redemption cache failed", logger.Field("error", err.Error()))
return err
}
} else {
// Redis miss — fall back to DB activation_context field.
logger.WithContext(ctx).Infow("Redis cache miss for redemption order, falling back to DB",
logger.Field("cache_key", cacheKey),
logger.Field("error", redisErr.Error()),
)
if orderInfo.ActivationContext == "" {
logger.WithContext(ctx).Error("CRITICAL: activation_context missing in DB and Redis expired; cannot recover redemption order",
logger.Field("order_no", orderInfo.OrderNo),
)
return errors.New("activation context not found: Redis expired and no DB fallback available")
}
var fullContext struct {
Type string `json:"type"`
RedemptionCodeId int64 `json:"redemption_code_id"`
UnitTime string `json:"unit_time"`
Quantity int64 `json:"quantity"`
}
if err = json.Unmarshal([]byte(orderInfo.ActivationContext), &fullContext); err != nil {
logger.WithContext(ctx).Error("Unmarshal DB activation_context failed",
logger.Field("order_no", orderInfo.OrderNo),
logger.Field("error", err.Error()),
)
return err
}
redemptionData.RedemptionCodeId = fullContext.RedemptionCodeId
redemptionData.UnitTime = fullContext.UnitTime
redemptionData.Quantity = fullContext.Quantity
} }
// 4. 幂等性检查:查询是否已有兑换记录 // 4. 幂等性检查:查询是否已有兑换记录
@@ -1,104 +0,0 @@
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
}
+2 -3
View File
@@ -5,7 +5,6 @@ const (
SchedulerTotalServerData = "scheduler:total:server" SchedulerTotalServerData = "scheduler:total:server"
SchedulerResetTraffic = "scheduler:reset:traffic" SchedulerResetTraffic = "scheduler:reset:traffic"
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 订单
) )
-6
View File
@@ -64,12 +64,6 @@ 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())
} }