Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 05ee5b9d2e |
+14
-14
@@ -38,20 +38,20 @@ type (
|
|||||||
Id int64 `form:"id" validate:"required"`
|
Id int64 `form:"id" validate:"required"`
|
||||||
}
|
}
|
||||||
UpdateUserBasiceInfoRequest {
|
UpdateUserBasiceInfoRequest {
|
||||||
UserId int64 `json:"user_id" validate:"required"`
|
UserId int64 `json:"user_id" validate:"required"`
|
||||||
Password string `json:"password"`
|
Password string `json:"password"`
|
||||||
Avatar string `json:"avatar"`
|
Avatar string `json:"avatar"`
|
||||||
Balance *int64 `json:"balance"`
|
Balance int64 `json:"balance"`
|
||||||
Commission *int64 `json:"commission"`
|
Commission int64 `json:"commission"`
|
||||||
ReferralPercentage uint8 `json:"referral_percentage"`
|
ReferralPercentage uint8 `json:"referral_percentage"`
|
||||||
OnlyFirstPurchase *bool `json:"only_first_purchase"`
|
OnlyFirstPurchase *bool `json:"only_first_purchase"`
|
||||||
GiftAmount *int64 `json:"gift_amount"`
|
GiftAmount int64 `json:"gift_amount"`
|
||||||
Telegram int64 `json:"telegram"`
|
Telegram int64 `json:"telegram"`
|
||||||
ReferCode string `json:"refer_code"`
|
ReferCode string `json:"refer_code"`
|
||||||
RefererId *int64 `json:"referer_id"`
|
RefererId int64 `json:"referer_id"`
|
||||||
Enable *bool `json:"enable"`
|
Enable *bool `json:"enable"`
|
||||||
IsAdmin *bool `json:"is_admin"`
|
IsAdmin *bool `json:"is_admin"`
|
||||||
Remark *string `json:"remark"`
|
Remark string `json:"remark"`
|
||||||
}
|
}
|
||||||
UpdateUserNotifySettingRequest {
|
UpdateUserNotifySettingRequest {
|
||||||
UserId int64 `json:"user_id" validate:"required"`
|
UserId int64 `json:"user_id" validate:"required"`
|
||||||
|
|||||||
@@ -46,13 +46,13 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
|||||||
}
|
}
|
||||||
|
|
||||||
err = l.svcCtx.UserModel.Transaction(l.ctx, func(tx *gorm.DB) error {
|
err = l.svcCtx.UserModel.Transaction(l.ctx, func(tx *gorm.DB) error {
|
||||||
if req.Balance != nil && userInfo.Balance != *req.Balance {
|
if userInfo.Balance != req.Balance {
|
||||||
change := *req.Balance - userInfo.Balance
|
change := req.Balance - userInfo.Balance
|
||||||
balanceLog := log.Balance{
|
balanceLog := log.Balance{
|
||||||
Type: log.BalanceTypeAdjust,
|
Type: log.BalanceTypeAdjust,
|
||||||
Amount: change,
|
Amount: change,
|
||||||
OrderNo: "",
|
OrderNo: "",
|
||||||
Balance: *req.Balance,
|
Balance: req.Balance,
|
||||||
Timestamp: time.Now().UnixMilli(),
|
Timestamp: time.Now().UnixMilli(),
|
||||||
}
|
}
|
||||||
content, _ := balanceLog.Marshal()
|
content, _ := balanceLog.Marshal()
|
||||||
@@ -66,14 +66,14 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
userInfo.Balance = *req.Balance
|
userInfo.Balance = req.Balance
|
||||||
}
|
}
|
||||||
|
|
||||||
if req.GiftAmount != nil && userInfo.GiftAmount != *req.GiftAmount {
|
if userInfo.GiftAmount != req.GiftAmount {
|
||||||
change := *req.GiftAmount - userInfo.GiftAmount
|
change := req.GiftAmount - userInfo.GiftAmount
|
||||||
if change != 0 {
|
if change != 0 {
|
||||||
var changeType uint16
|
var changeType uint16
|
||||||
if userInfo.GiftAmount < *req.GiftAmount {
|
if userInfo.GiftAmount < req.GiftAmount {
|
||||||
changeType = log.GiftTypeIncrease
|
changeType = log.GiftTypeIncrease
|
||||||
} else {
|
} else {
|
||||||
changeType = log.GiftTypeReduce
|
changeType = log.GiftTypeReduce
|
||||||
@@ -81,7 +81,7 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
|||||||
giftLog := log.Gift{
|
giftLog := log.Gift{
|
||||||
Type: changeType,
|
Type: changeType,
|
||||||
Amount: change,
|
Amount: change,
|
||||||
Balance: *req.GiftAmount,
|
Balance: req.GiftAmount,
|
||||||
Remark: "Admin adjustment",
|
Remark: "Admin adjustment",
|
||||||
Timestamp: time.Now().UnixMilli(),
|
Timestamp: time.Now().UnixMilli(),
|
||||||
}
|
}
|
||||||
@@ -96,27 +96,23 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
userInfo.GiftAmount = *req.GiftAmount
|
userInfo.GiftAmount = req.GiftAmount
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if req.Commission != nil && *req.Commission != userInfo.Commission {
|
if req.Commission != userInfo.Commission {
|
||||||
remark := ""
|
if isWithdrawalScene(req.Remark) {
|
||||||
if req.Remark != nil {
|
|
||||||
remark = *req.Remark
|
|
||||||
}
|
|
||||||
if isWithdrawalScene(remark) {
|
|
||||||
logWithdrawalGuard(l.Logger, userInfo.Id)
|
logWithdrawalGuard(l.Logger, userInfo.Id)
|
||||||
return errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "commission overwrite is blocked in withdrawal scene")
|
return errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "commission overwrite is blocked in withdrawal scene")
|
||||||
}
|
}
|
||||||
change := *req.Commission - userInfo.Commission
|
change := req.Commission - userInfo.Commission
|
||||||
if err = l.svcCtx.UserModel.UpdateCommission(l.ctx, userInfo.Id, change, tx); err != nil {
|
if err = l.svcCtx.UserModel.UpdateCommission(l.ctx, userInfo.Id, change, tx); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err = logicCommon.WriteCommissionLog(tx, userInfo.Id, log.CommissionTypeAdjust, change, ""); err != nil {
|
if err = logicCommon.WriteCommissionLog(tx, userInfo.Id, log.CommissionTypeAdjust, change, ""); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
userInfo.Commission = *req.Commission
|
userInfo.Commission = req.Commission
|
||||||
}
|
}
|
||||||
if req.Avatar != "" {
|
if req.Avatar != "" {
|
||||||
userInfo.Avatar = req.Avatar
|
userInfo.Avatar = req.Avatar
|
||||||
@@ -124,17 +120,15 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
|||||||
if req.ReferCode != "" {
|
if req.ReferCode != "" {
|
||||||
userInfo.ReferCode = req.ReferCode
|
userInfo.ReferCode = req.ReferCode
|
||||||
}
|
}
|
||||||
if req.RefererId != nil {
|
userInfo.RefererId = req.RefererId
|
||||||
userInfo.RefererId = *req.RefererId
|
|
||||||
}
|
|
||||||
if req.Enable != nil {
|
if req.Enable != nil {
|
||||||
userInfo.Enable = req.Enable
|
userInfo.Enable = req.Enable
|
||||||
}
|
}
|
||||||
if req.IsAdmin != nil {
|
if req.IsAdmin != nil {
|
||||||
userInfo.IsAdmin = req.IsAdmin
|
userInfo.IsAdmin = req.IsAdmin
|
||||||
}
|
}
|
||||||
if req.Remark != nil {
|
if req.Remark != "" {
|
||||||
userInfo.Remark = *req.Remark
|
userInfo.Remark = req.Remark
|
||||||
}
|
}
|
||||||
if req.OnlyFirstPurchase != nil {
|
if req.OnlyFirstPurchase != nil {
|
||||||
userInfo.OnlyFirstPurchase = req.OnlyFirstPurchase
|
userInfo.OnlyFirstPurchase = req.OnlyFirstPurchase
|
||||||
|
|||||||
@@ -3,17 +3,20 @@ 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"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -84,6 +87,7 @@ 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
|
||||||
@@ -96,6 +100,32 @@ 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)
|
||||||
@@ -169,7 +199,12 @@ func (l *AppleIAPNotifyLogic) Handle(signedPayload string) error {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
if days == 0 {
|
if days == 0 {
|
||||||
l.Errorw("iap notify product mapping missing", logger.Field("productId", txPayload.ProductId))
|
// Both string-parse and DB fallback failed to map the product to days.
|
||||||
|
// 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)
|
||||||
@@ -216,6 +251,11 @@ 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
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -248,7 +288,52 @@ 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
|
||||||
|
}
|
||||||
|
|||||||
@@ -110,6 +110,10 @@ 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]
|
||||||
@@ -123,7 +127,7 @@ func (m *customOrderModel) UpdateOrderStatus(ctx context.Context, orderNo string
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}, m.getCacheKeys(orderInfo)...)
|
}, keys...)
|
||||||
}
|
}
|
||||||
|
|
||||||
// FindOneDetailsByOrderNo Find order details by order number
|
// FindOneDetailsByOrderNo Find order details by order number
|
||||||
|
|||||||
+14
-14
@@ -3223,20 +3223,20 @@ type UpdateUserAuthMethodRequest struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type UpdateUserBasiceInfoRequest struct {
|
type UpdateUserBasiceInfoRequest struct {
|
||||||
UserId int64 `json:"user_id" validate:"required"`
|
UserId int64 `json:"user_id" validate:"required"`
|
||||||
Password string `json:"password"`
|
Password string `json:"password"`
|
||||||
Avatar string `json:"avatar"`
|
Avatar string `json:"avatar"`
|
||||||
Balance *int64 `json:"balance"`
|
Balance int64 `json:"balance"`
|
||||||
Commission *int64 `json:"commission"`
|
Commission int64 `json:"commission"`
|
||||||
ReferralPercentage uint8 `json:"referral_percentage"`
|
ReferralPercentage uint8 `json:"referral_percentage"`
|
||||||
OnlyFirstPurchase *bool `json:"only_first_purchase"`
|
OnlyFirstPurchase *bool `json:"only_first_purchase"`
|
||||||
GiftAmount *int64 `json:"gift_amount"`
|
GiftAmount int64 `json:"gift_amount"`
|
||||||
Telegram int64 `json:"telegram"`
|
Telegram int64 `json:"telegram"`
|
||||||
ReferCode string `json:"refer_code"`
|
ReferCode string `json:"refer_code"`
|
||||||
RefererId *int64 `json:"referer_id"`
|
RefererId int64 `json:"referer_id"`
|
||||||
Enable *bool `json:"enable"`
|
Enable *bool `json:"enable"`
|
||||||
IsAdmin *bool `json:"is_admin"`
|
IsAdmin *bool `json:"is_admin"`
|
||||||
Remark *string `json:"remark"`
|
Remark string `json:"remark"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type UpdateUserNotifyRequest struct {
|
type UpdateUserNotifyRequest struct {
|
||||||
|
|||||||
@@ -1,33 +0,0 @@
|
|||||||
package types
|
|
||||||
|
|
||||||
import (
|
|
||||||
"encoding/json"
|
|
||||||
"testing"
|
|
||||||
)
|
|
||||||
|
|
||||||
func TestUpdateUserBasiceInfoRequestDistinguishesOmittedAndEmptyRemark(t *testing.T) {
|
|
||||||
var omitted UpdateUserBasiceInfoRequest
|
|
||||||
if err := json.Unmarshal([]byte(`{"user_id":1001}`), &omitted); err != nil {
|
|
||||||
t.Fatalf("unmarshal omitted remark: %v", err)
|
|
||||||
}
|
|
||||||
if omitted.Remark != nil {
|
|
||||||
t.Fatalf("omitted remark should stay nil, got %q", *omitted.Remark)
|
|
||||||
}
|
|
||||||
if omitted.RefererId != nil {
|
|
||||||
t.Fatalf("omitted referer_id should stay nil, got %d", *omitted.RefererId)
|
|
||||||
}
|
|
||||||
if omitted.Balance != nil || omitted.GiftAmount != nil || omitted.Commission != nil {
|
|
||||||
t.Fatalf("omitted money fields should stay nil, got balance=%v gift=%v commission=%v", omitted.Balance, omitted.GiftAmount, omitted.Commission)
|
|
||||||
}
|
|
||||||
|
|
||||||
var cleared UpdateUserBasiceInfoRequest
|
|
||||||
if err := json.Unmarshal([]byte(`{"user_id":1001,"remark":""}`), &cleared); err != nil {
|
|
||||||
t.Fatalf("unmarshal empty remark: %v", err)
|
|
||||||
}
|
|
||||||
if cleared.Remark == nil {
|
|
||||||
t.Fatal("explicit empty remark should be present")
|
|
||||||
}
|
|
||||||
if *cleared.Remark != "" {
|
|
||||||
t.Fatalf("explicit empty remark should decode to empty string, got %q", *cleared.Remark)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -597,7 +597,7 @@ func (l *ActivateOrderLogic) NewPurchase(ctx context.Context, orderInfo *order.O
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
if err = validateNewUserOnlyEligibilityAtActivation(ctx, l.svc.DB, orderInfo, sub); err != nil {
|
if err = validateNewUserOnlyEligibilityAtActivation(ctx, l.svc.DB, l.svc.Redis, orderInfo, sub); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -675,18 +675,23 @@ 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
|
||||||
if findErr := l.svc.DB.Model(&user.Subscribe{}).
|
findErr := l.svc.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||||
Where("user_id IN ? AND token != ''", candidateUserIds).
|
return tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||||
Order("expire_time DESC").
|
Model(&user.Subscribe{}).
|
||||||
Order("updated_at DESC").
|
Where("user_id IN ? AND token != ''", candidateUserIds).
|
||||||
Order("id DESC").
|
Order("expire_time DESC").
|
||||||
First(&existingSub).Error; findErr == nil {
|
Order("updated_at DESC").
|
||||||
|
Order("id DESC").
|
||||||
|
First(&existingSub).Error
|
||||||
|
})
|
||||||
|
if findErr == nil {
|
||||||
// 家庭组场景:订阅 owner 可能变更(如成员注册的试用 → 被家主收归),
|
// 家庭组场景:订阅 owner 可能变更(如成员注册的试用 → 被家主收归),
|
||||||
// 续期前把 user_id 校正为当前订单的 SubscriptionUserId
|
// 续期前把 user_id 校正为当前订单的 SubscriptionUserId
|
||||||
effectiveOwner := orderInfo.UserId
|
effectiveOwner := orderInfo.UserId
|
||||||
@@ -1447,7 +1452,40 @@ 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 {
|
||||||
return nil, err
|
// Fallback: token may have been lost (subscription deleted or owner changed).
|
||||||
|
// 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
|
||||||
|
|||||||
@@ -10,12 +10,14 @@ 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 {
|
||||||
@@ -31,6 +33,20 @@ 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
|
||||||
|
|||||||
Reference in New Issue
Block a user