Compare commits

..

1 Commits

Author SHA1 Message Date
shanshanzhong147 af22430101 fix: 修复订单 claim 机制状态冲突与 stuck 订单恢复
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>
2026-05-24 23:42:37 -07:00
12 changed files with 114 additions and 341 deletions
@@ -1,2 +0,0 @@
-- Remove activation_context column from order table
ALTER TABLE `order` DROP COLUMN IF EXISTS `activation_context`;
@@ -1,6 +0,0 @@
-- 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;
+1 -86
View File
@@ -3,20 +3,17 @@ package notify
import ( import (
"context" "context"
"encoding/json" "encoding/json"
"fmt"
"strconv" "strconv"
"strings" "strings"
commonLogic "github.com/perfect-panel/server/internal/logic/common" commonLogic "github.com/perfect-panel/server/internal/logic/common"
iapmodel "github.com/perfect-panel/server/internal/model/iap/apple" iapmodel "github.com/perfect-panel/server/internal/model/iap/apple"
"github.com/perfect-panel/server/internal/model/order"
"github.com/perfect-panel/server/internal/model/subscribe" "github.com/perfect-panel/server/internal/model/subscribe"
"github.com/perfect-panel/server/internal/model/user" "github.com/perfect-panel/server/internal/model/user"
"github.com/perfect-panel/server/internal/svc" "github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types" "github.com/perfect-panel/server/internal/types"
iapapple "github.com/perfect-panel/server/pkg/iap/apple" iapapple "github.com/perfect-panel/server/pkg/iap/apple"
"github.com/perfect-panel/server/pkg/logger" "github.com/perfect-panel/server/pkg/logger"
"github.com/perfect-panel/server/pkg/tool"
"gorm.io/gorm" "gorm.io/gorm"
) )
@@ -87,7 +84,6 @@ func (l *AppleIAPNotifyLogic) Handle(signedPayload string) error {
l.Errorw("iap notify insert transaction error", logger.Field("error", e.Error()), logger.Field("productId", txPayload.ProductId), logger.Field("originalTransactionId", txPayload.OriginalTransactionId)) l.Errorw("iap notify insert transaction error", logger.Field("error", e.Error()), logger.Field("productId", txPayload.ProductId), logger.Field("originalTransactionId", txPayload.OriginalTransactionId))
return e return e
} }
existing = rec
} else { } else {
if txPayload.RevocationDate != nil { if txPayload.RevocationDate != nil {
// 撤销场景:更新 revocation_at // 撤销场景:更新 revocation_at
@@ -100,32 +96,6 @@ func (l *AppleIAPNotifyLogic) Handle(signedPayload string) error {
} }
} }
} }
// Fix Bug 7: resolve userId when the IAP transaction record has no associated user.
// This can happen if the first notification for a subscription arrived before the
// in-app purchase flow created the order (race) or the order was made in a previous
// build that did not write user_id to the IAP transaction table.
if existing != nil && existing.UserId == 0 {
var origOrder order.Order
if lookupErr := db.Model(&order.Order{}).
Where("trade_no = ? AND method = ?", txPayload.OriginalTransactionId, "apple_iap").
Order("id ASC").
First(&origOrder).Error; lookupErr == nil && origOrder.UserId > 0 {
existing.UserId = origOrder.UserId
_ = db.Model(&iapmodel.Transaction{}).
Where("id = ?", existing.Id).
Update("user_id", origOrder.UserId).Error
l.Infow("iap notify resolved zero userId from original purchase order",
logger.Field("userId", origOrder.UserId),
logger.Field("originalTransactionId", txPayload.OriginalTransactionId),
)
} else {
l.Errorw("CRITICAL: iap notify UserId=0 and cannot resolve from order, notification dropped",
logger.Field("originalTransactionId", txPayload.OriginalTransactionId))
return fmt.Errorf("iap notify: UserId=0 and cannot resolve from order for original_transaction_id=%s", txPayload.OriginalTransactionId)
}
}
var days int64 var days int64
{ {
pid := strings.ToLower(txPayload.ProductId) pid := strings.ToLower(txPayload.ProductId)
@@ -199,12 +169,7 @@ func (l *AppleIAPNotifyLogic) Handle(signedPayload string) error {
} }
} }
if days == 0 { if days == 0 {
// Both string-parse and DB fallback failed to map the product to days. l.Errorw("iap notify product mapping missing", logger.Field("productId", txPayload.ProductId))
// Return an error so Apple retries the notification; silently ignoring
// this would cause the subscriber's renewal to be lost.
l.Errorw("CRITICAL: iap notify product mapping missing, returning error to trigger retry",
logger.Field("productId", txPayload.ProductId))
return fmt.Errorf("iap product id %s could not be mapped to subscription days", txPayload.ProductId)
} }
token := "iap:" + txPayload.OriginalTransactionId token := "iap:" + txPayload.OriginalTransactionId
sub, e := l.svcCtx.UserModel.FindOneSubscribeByToken(l.ctx, token) sub, e := l.svcCtx.UserModel.FindOneSubscribeByToken(l.ctx, token)
@@ -251,11 +216,6 @@ func (l *AppleIAPNotifyLogic) Handle(signedPayload string) error {
logger.Field("product_id", txPayload.ProductId), logger.Field("product_id", txPayload.ProductId),
)..., )...,
) )
// Create audit order record for renewal notifications (Bug 7)
if err := l.createIAPRenewalAuditOrder(db, ntype, txPayload.TransactionId, candidate.UserId, candidate.SubscribeId, candidate.Token); err != nil {
l.Errorw("iap notify fallback create renewal order error", logger.Field("error", err.Error()))
// Non-fatal: subscription already updated; order creation failure is logged only
}
break break
} }
} }
@@ -288,52 +248,7 @@ func (l *AppleIAPNotifyLogic) Handle(signedPayload string) error {
logger.Field("product_id", txPayload.ProductId), logger.Field("product_id", txPayload.ProductId),
)..., )...,
) )
// Create audit order record for renewal notifications (Bug 7)
if err := l.createIAPRenewalAuditOrder(db, ntype, txPayload.TransactionId, sub.UserId, sub.SubscribeId, sub.Token); err != nil {
l.Errorw("iap notify create renewal order error", logger.Field("error", err.Error()))
// Non-fatal: subscription already updated; order creation failure is logged only
}
} }
return nil return nil
}) })
} }
// createIAPRenewalAuditOrder creates a finished renewal order record for DID_RENEW and SUBSCRIBED
// Apple SSNS notifications so that the order table has a complete financial audit trail.
// The operation is idempotent: if an order with the same trade_no already exists it is skipped.
func (l *AppleIAPNotifyLogic) createIAPRenewalAuditOrder(db *gorm.DB, ntype, transactionId string, userId, subscribeId int64, subscribeToken string) error {
if ntype != "DID_RENEW" && ntype != "SUBSCRIBED" {
return nil
}
// Idempotency check
var count int64
if err := db.Model(&order.Order{}).
Where("trade_no = ? AND method = ?", transactionId, "apple_iap").
Count(&count).Error; err != nil {
return fmt.Errorf("check existing iap renewal order: %w", err)
}
if count > 0 {
return nil
}
rec := &order.Order{
UserId: userId,
OrderNo: tool.GenerateTradeNo(),
Type: 2, // OrderTypeRenewal
Status: 5, // OrderStatusFinished
Method: "apple_iap",
TradeNo: transactionId,
SubscribeId: subscribeId,
SubscribeToken: subscribeToken,
Quantity: 1,
IsNew: false,
}
if err := db.Model(&order.Order{}).Create(rec).Error; err != nil {
return fmt.Errorf("create iap renewal audit order: %w", err)
}
l.Infow("iap notify created renewal audit order",
logger.Field("orderNo", rec.OrderNo),
logger.Field("transactionId", transactionId),
logger.Field("userId", userId),
)
return nil
}
@@ -154,12 +154,9 @@ func (l *PurchaseLogic) Purchase(req *types.PortalPurchaseRequest) (resp *types.
} }
content, _ := tempOrder.Marshal() content, _ := tempOrder.Marshal()
// Persist activation context to DB so the worker can recover if Redis TTL expires. if _, err = l.svcCtx.Redis.Set(l.ctx, fmt.Sprintf(constant.TempOrderCacheKey, orderInfo.OrderNo), string(content), CloseOrderTimeMinutes*time.Minute).Result(); err != nil {
orderInfo.ActivationContext = string(content) l.Errorw("[Purchase] Redis set error", logger.Field("error", err.Error()), logger.Field("order_no", orderInfo.OrderNo))
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))
@@ -172,7 +169,7 @@ func (l *PurchaseLogic) Purchase(req *types.PortalPurchaseRequest) (resp *types.
} }
} }
// save guest order (activation_context is included) // save guest order
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,48 +151,34 @@ 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到数据库activation_context 同步写入,作为 Redis 的持久化兜底) // 保存Order到数据库
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,
@@ -200,8 +186,16 @@ func (l *RedeemCodeLogic) RedeemCode(req *types.RedeemCodeRequest) (resp *types.
"quantity": redemptionCode.Quantity, "quantity": redemptionCode.Quantity,
} }
jsonData, _ := json.Marshal(cacheData) jsonData, _ := json.Marshal(cacheData)
if redisErr := l.svcCtx.Redis.Set(l.ctx, cacheKey, jsonData, 2*time.Hour).Err(); redisErr != nil { err = l.svcCtx.Redis.Set(l.ctx, cacheKey, jsonData, 2*time.Hour).Err()
l.Infow("[RedeemCode] Cache redemption data failed (non-fatal, DB fallback available)", logger.Field("error", redisErr.Error())) if err != nil {
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")
} }
// 触发队列任务 // 触发队列任务
+1 -5
View File
@@ -110,10 +110,6 @@ func (m *customOrderModel) UpdateOrderStatus(ctx context.Context, orderNo string
if err != nil { if err != nil {
return err return err
} }
keys := m.getCacheKeys(orderInfo)
// Pre-delete: evict cache before the DB write so concurrent reads during the update
// window go to DB instead of getting a stale cached status (double-delete pattern).
_ = m.DelCacheCtx(ctx, keys...)
return m.ExecCtx(ctx, func(conn *gorm.DB) error { return m.ExecCtx(ctx, func(conn *gorm.DB) error {
if len(tx) > 0 { if len(tx) > 0 {
conn = tx[0] conn = tx[0]
@@ -127,7 +123,7 @@ func (m *customOrderModel) UpdateOrderStatus(ctx context.Context, orderNo string
return nil return nil
} }
return nil return nil
}, keys...) }, m.getCacheKeys(orderInfo)...)
} }
// FindOneDetailsByOrderNo Find order details by order number // FindOneDetailsByOrderNo Find order details by order number
+2 -3
View File
@@ -24,9 +24,8 @@ 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)"`
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"`
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"`
} }
+1 -1
View File
@@ -49,6 +49,6 @@ func RegisterHandlers(mux *asynq.ServeMux, serverCtx *svc.ServiceContext) {
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 // Stuck order recovery: reset claimed orders that timed out back to paid
mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx)) mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx))
} }
+47 -133
View File
@@ -93,11 +93,12 @@ 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") { if strings.Contains(err.Error(), "stuck in claimed") {
logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单卡在 claimed,将重试", logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单卡在 claimed,将重试",
logger.Field("order_no", payload.OrderNo)) logger.Field("order_no", payload.OrderNo))
return err // 返回错误触发 asynq 重试 return err
} }
logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过", logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过",
logger.Field("order_no", payload.OrderNo)) logger.Field("order_no", payload.OrderNo))
@@ -191,7 +192,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 { if orderInfo.Status == OrderStatusClaimed {
logger.WithContext(ctx).Error("Order stuck in claimed status", logger.WithContext(ctx).Error("Order stuck in claimed status",
logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_no", orderInfo.OrderNo),
@@ -589,9 +589,8 @@ func (l *ActivateOrderLogic) finalizeCouponAndOrder(ctx context.Context, orderIn
} }
} }
// UpdateOrderStatus uses WHERE status < target, which blocks claimed(6)finished(5). // Direct update from claimed(6)finished(5), bypassing the model's status<target guard.
// Use a direct update matching the exact claimed status, then update the full record // UpdateOrderStatus uses WHERE status < ?, so status=6 > 5 would never match.
// via the model layer to properly invalidate the cache.
result := l.svc.DB.WithContext(ctx). result := l.svc.DB.WithContext(ctx).
Model(&order.Order{}). Model(&order.Order{}).
Where("order_no = ? AND status = ?", orderInfo.OrderNo, OrderStatusClaimed). Where("order_no = ? AND status = ?", orderInfo.OrderNo, OrderStatusClaimed).
@@ -602,14 +601,18 @@ func (l *ActivateOrderLogic) finalizeCouponAndOrder(ctx context.Context, orderIn
logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_no", orderInfo.OrderNo),
) )
} }
// Invalidate order cache regardless of whether the DB update succeeded // Invalidate cache entries; key format matches order/default.go cacheOrderIdPrefix / cacheOrderNoPrefix.
orderInfo.Status = OrderStatusFinished cacheKeys := []string{
if err := l.svc.OrderModel.Update(ctx, orderInfo); err != nil { fmt.Sprintf("cache:order:id:%d", orderInfo.Id),
logger.WithContext(ctx).Error("Update order cache after finalization failed", fmt.Sprintf("cache:order:no:%s", orderInfo.OrderNo),
logger.Field("error", err.Error()), }
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),
) )
} }
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)...,
@@ -636,7 +639,7 @@ func (l *ActivateOrderLogic) NewPurchase(ctx context.Context, orderInfo *order.O
return err return err
} }
if err = validateNewUserOnlyEligibilityAtActivation(ctx, l.svc.DB, l.svc.Redis, orderInfo, sub); err != nil { if err = validateNewUserOnlyEligibilityAtActivation(ctx, l.svc.DB, orderInfo, sub); err != nil {
return err return err
} }
@@ -714,23 +717,18 @@ func (l *ActivateOrderLogic) NewPurchase(ctx context.Context, orderInfo *order.O
// 兜底:创建新订阅前,查找用户是否已有同套餐的订阅记录(含过期/赠送), // 兜底:创建新订阅前,查找用户是否已有同套餐的订阅记录(含过期/赠送),
// 有则复用旧记录续期,避免出现重复订阅。 // 有则复用旧记录续期,避免出现重复订阅。
// 需要同时检查 UserId 和 SubscriptionUserId,因为家庭组绑定前后 owner 可能不同。 // 需要同时检查 UserId 和 SubscriptionUserId,因为家庭组绑定前后 owner 可能不同。
// 使用 SELECT ... FOR UPDATE 防止并发下多个 worker 同时选中同一订阅并创建重复记录。
if userSub == nil { if userSub == nil {
candidateUserIds := []int64{orderInfo.UserId} candidateUserIds := []int64{orderInfo.UserId}
if orderInfo.SubscriptionUserId > 0 && orderInfo.SubscriptionUserId != orderInfo.UserId { if orderInfo.SubscriptionUserId > 0 && orderInfo.SubscriptionUserId != orderInfo.UserId {
candidateUserIds = append(candidateUserIds, orderInfo.SubscriptionUserId) candidateUserIds = append(candidateUserIds, orderInfo.SubscriptionUserId)
} }
var existingSub user.Subscribe var existingSub user.Subscribe
findErr := l.svc.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { if findErr := l.svc.DB.Model(&user.Subscribe{}).
return tx.Clauses(clause.Locking{Strength: "UPDATE"}). Where("user_id IN ? AND token != ''", candidateUserIds).
Model(&user.Subscribe{}). Order("expire_time DESC").
Where("user_id IN ? AND token != ''", candidateUserIds). Order("updated_at DESC").
Order("expire_time DESC"). Order("id DESC").
Order("updated_at DESC"). First(&existingSub).Error; findErr == nil {
Order("id DESC").
First(&existingSub).Error
})
if findErr == nil {
// 家庭组场景:订阅 owner 可能变更(如成员注册的试用 → 被家主收归), // 家庭组场景:订阅 owner 可能变更(如成员注册的试用 → 被家主收归),
// 续期前把 user_id 校正为当前订单的 SubscriptionUserId // 续期前把 user_id 校正为当前订单的 SubscriptionUserId
effectiveOwner := orderInfo.UserId effectiveOwner := orderInfo.UserId
@@ -868,52 +866,26 @@ func (l *ActivateOrderLogic) createGuestUser(ctx context.Context, orderInfo *ord
return userInfo, nil return userInfo, nil
} }
// getTempOrderInfo retrieves temporary order information from Redis cache with DB fallback. // getTempOrderInfo retrieves temporary order information from Redis cache
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 {
var tempOrder constant.TemporaryOrderInfo logger.WithContext(ctx).Error("Get temp order cache failed",
if unmarshalErr := tempOrder.Unmarshal([]byte(data)); unmarshalErr != nil { logger.Field("error", err.Error()),
logger.WithContext(ctx).Error("Unmarshal temp order cache failed", logger.Field("cache_key", cacheKey),
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, dbErr return nil, err
}
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 unmarshalErr := tempOrder.Unmarshal([]byte(orderInfo.ActivationContext)); unmarshalErr != nil { if err = tempOrder.Unmarshal([]byte(data)); err != nil {
logger.WithContext(ctx).Error("Unmarshal DB activation_context failed", logger.WithContext(ctx).Error("Unmarshal temp order cache failed",
logger.Field("order_no", orderNo), logger.Field("error", err.Error()),
logger.Field("error", unmarshalErr.Error()), logger.Field("cache_key", cacheKey),
logger.Field("data", data),
) )
return nil, unmarshalErr return nil, err
} }
return &tempOrder, nil return &tempOrder, nil
@@ -1517,40 +1489,7 @@ func (l *ActivateOrderLogic) resolveRenewalActivationSubscription(ctx context.Co
} }
userSub, err := l.getUserSubscription(ctx, orderInfo.SubscribeToken) userSub, err := l.getUserSubscription(ctx, orderInfo.SubscribeToken)
if err != nil { if err != nil {
// Fallback: token may have been lost (subscription deleted or owner changed). return nil, err
// Locate the most recent subscription by user_id + subscribe_id with FOR UPDATE
// to prevent a concurrent worker from renewing the same record twice.
targetUserID := orderInfo.UserId
if orderInfo.SubscriptionUserId > 0 {
targetUserID = orderInfo.SubscriptionUserId
}
var fallbackSub user.Subscribe
txErr := l.svc.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
return tx.Clauses(clause.Locking{Strength: "UPDATE"}).
Model(&user.Subscribe{}).
Where("user_id = ? AND subscribe_id = ?", targetUserID, orderInfo.SubscribeId).
Where("status IN ?", []int64{0, 1, 2, 3}).
Order("expire_time DESC").
Order("updated_at DESC").
Order("id DESC").
First(&fallbackSub).Error
})
if txErr != nil {
logger.WithContext(ctx).Error("CRITICAL: Renewal activation token and fallback lookup both failed",
logger.Field("token_error", err.Error()),
logger.Field("fallback_error", txErr.Error()),
logger.Field("order_no", orderInfo.OrderNo),
logger.Field("target_user_id", targetUserID),
logger.Field("subscribe_id", orderInfo.SubscribeId),
)
return nil, fmt.Errorf("renewal activation subscription not found by token or user_id+subscribe_id for order %s: %w", orderInfo.OrderNo, err)
}
logger.WithContext(ctx).Info("Renewal token lookup failed; found subscription via fallback user_id+subscribe_id",
logger.Field("order_no", orderInfo.OrderNo),
logger.Field("fallback_subscribe_id", fallbackSub.Id),
logger.Field("target_user_id", targetUserID),
)
userSub = &fallbackSub
} }
if orderInfo.UserId <= 0 { if orderInfo.UserId <= 0 {
return userSub, nil return userSub, nil
@@ -1846,50 +1785,25 @@ func (l *ActivateOrderLogic) RedemptionActivate(ctx context.Context, orderInfo *
return err return err
} }
// 3. 从Redis获取兑换码信息Redis缺失时回查DB // 3. 从Redis获取兑换码信息
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 {
cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo) logger.WithContext(ctx).Error("Unmarshal redemption cache failed", logger.Field("error", err.Error()))
data, redisErr := l.svc.Redis.Get(ctx, cacheKey).Result() return err
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. 幂等性检查:查询是否已有兑换记录
-16
View File
@@ -10,14 +10,12 @@ import (
"github.com/perfect-panel/server/internal/model/order" "github.com/perfect-panel/server/internal/model/order"
"github.com/perfect-panel/server/internal/model/subscribe" "github.com/perfect-panel/server/internal/model/subscribe"
internaltypes "github.com/perfect-panel/server/internal/types" internaltypes "github.com/perfect-panel/server/internal/types"
"github.com/redis/go-redis/v9"
"gorm.io/gorm" "gorm.io/gorm"
) )
func validateNewUserOnlyEligibilityAtActivation( func validateNewUserOnlyEligibilityAtActivation(
ctx context.Context, ctx context.Context,
db *gorm.DB, db *gorm.DB,
rdb *redis.Client,
orderInfo *order.Order, orderInfo *order.Order,
sub *subscribe.Subscribe, sub *subscribe.Subscribe,
) error { ) error {
@@ -33,20 +31,6 @@ func validateNewUserOnlyEligibilityAtActivation(
return nil return nil
} }
// Acquire a per-user distributed lock so concurrent new-user-only activations
// for the same account are serialised. Without this, two workers can both read
// historyCount=0 and both pass the check before either has written the order.
lockKey := fmt.Sprintf("new_user_only_activate:%d", orderInfo.UserId)
const lockTTL = 30 * time.Second
acquired, lockErr := rdb.SetNX(ctx, lockKey, orderInfo.OrderNo, lockTTL).Result()
if lockErr != nil {
return fmt.Errorf("new user only: acquire lock error: %w", lockErr)
}
if !acquired {
return fmt.Errorf("new user only: another activation is in progress for user %d", orderInfo.UserId)
}
defer rdb.Del(ctx, lockKey)
eligibility, err := commonLogic.ResolveNewUserEligibility(ctx, db, orderInfo.UserId) eligibility, err := commonLogic.ResolveNewUserEligibility(ctx, db, orderInfo.UserId)
if err != nil { if err != nil {
return err return err
+26 -44
View File
@@ -12,7 +12,10 @@ import (
queueTypes "github.com/perfect-panel/server/queue/types" queueTypes "github.com/perfect-panel/server/queue/types"
) )
// StuckOrderRecoveryLogic scans orders stuck in claimed status and re-queues them for processing. 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 { type StuckOrderRecoveryLogic struct {
svc *svc.ServiceContext svc *svc.ServiceContext
} }
@@ -21,11 +24,8 @@ func NewStuckOrderRecoveryLogic(svc *svc.ServiceContext) *StuckOrderRecoveryLogi
return &StuckOrderRecoveryLogic{svc: svc} 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 { func (l *StuckOrderRecoveryLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error {
cutoff := time.Now().Add(-10 * time.Minute) cutoff := time.Now().Add(-stuckClaimTimeout)
var stuckOrders []order.Order var stuckOrders []order.Order
if err := l.svc.DB.WithContext(ctx). if err := l.svc.DB.WithContext(ctx).
@@ -46,57 +46,39 @@ func (l *StuckOrderRecoveryLogic) ProcessTask(ctx context.Context, _ *asynq.Task
for i := range stuckOrders { for i := range stuckOrders {
orderNos = append(orderNos, stuckOrders[i].OrderNo) orderNos = append(orderNos, stuckOrders[i].OrderNo)
} }
logger.WithContext(ctx).Error("[StuckOrderRecovery] Found stuck claimed orders, recovering",
logger.WithContext(ctx).Error("[StuckOrderRecovery] Found stuck claimed orders, resetting to paid",
logger.Field("count", len(stuckOrders)), logger.Field("count", len(stuckOrders)),
logger.Field("order_nos", orderNos), 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 { for i := range stuckOrders {
o := &stuckOrders[i] ord := &stuckOrders[i]
payload, err := json.Marshal(queueTypes.ForthwithActivateOrderPayload{OrderNo: ord.OrderNo})
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 { if err != nil {
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to marshal task payload", logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to marshal payload",
logger.Field("order_no", o.OrderNo), logger.Field("order_no", ord.OrderNo),
logger.Field("error", err.Error()), logger.Field("error", err.Error()),
) )
continue continue
} }
if _, err = l.svc.Queue.EnqueueContext(ctx, asynq.NewTask(queueTypes.ForthwithActivateOrder, payload)); err != nil { task := asynq.NewTask(queueTypes.ForthwithActivateOrder, payload, asynq.MaxRetry(5))
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to re-enqueue activate task", if _, err := l.svc.Queue.EnqueueContext(ctx, task); err != nil {
logger.Field("order_no", o.OrderNo), logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to re-enqueue order",
logger.Field("order_no", ord.OrderNo),
logger.Field("error", err.Error()), logger.Field("error", err.Error()),
) )
} else {
logger.WithContext(ctx).Info("[StuckOrderRecovery] Re-enqueued activate task",
logger.Field("order_no", o.OrderNo),
)
} }
} }
+3 -3
View File
@@ -5,7 +5,7 @@ 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 订单 SchedulerStuckOrderRecovery = "scheduler:stuck:order:recovery" // 定时恢复卡在 claimed 状态的订单
) )