Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| af22430101 |
+5
-5
@@ -41,17 +41,17 @@ type (
|
||||
UserId int64 `json:"user_id" validate:"required"`
|
||||
Password string `json:"password"`
|
||||
Avatar string `json:"avatar"`
|
||||
Balance *int64 `json:"balance"`
|
||||
Commission *int64 `json:"commission"`
|
||||
Balance int64 `json:"balance"`
|
||||
Commission int64 `json:"commission"`
|
||||
ReferralPercentage uint8 `json:"referral_percentage"`
|
||||
OnlyFirstPurchase *bool `json:"only_first_purchase"`
|
||||
GiftAmount *int64 `json:"gift_amount"`
|
||||
GiftAmount int64 `json:"gift_amount"`
|
||||
Telegram int64 `json:"telegram"`
|
||||
ReferCode string `json:"refer_code"`
|
||||
RefererId *int64 `json:"referer_id"`
|
||||
RefererId int64 `json:"referer_id"`
|
||||
Enable *bool `json:"enable"`
|
||||
IsAdmin *bool `json:"is_admin"`
|
||||
Remark *string `json:"remark"`
|
||||
Remark string `json:"remark"`
|
||||
}
|
||||
UpdateUserNotifySettingRequest {
|
||||
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 {
|
||||
if req.Balance != nil && userInfo.Balance != *req.Balance {
|
||||
change := *req.Balance - userInfo.Balance
|
||||
if userInfo.Balance != req.Balance {
|
||||
change := req.Balance - userInfo.Balance
|
||||
balanceLog := log.Balance{
|
||||
Type: log.BalanceTypeAdjust,
|
||||
Amount: change,
|
||||
OrderNo: "",
|
||||
Balance: *req.Balance,
|
||||
Balance: req.Balance,
|
||||
Timestamp: time.Now().UnixMilli(),
|
||||
}
|
||||
content, _ := balanceLog.Marshal()
|
||||
@@ -66,14 +66,14 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
userInfo.Balance = *req.Balance
|
||||
userInfo.Balance = req.Balance
|
||||
}
|
||||
|
||||
if req.GiftAmount != nil && userInfo.GiftAmount != *req.GiftAmount {
|
||||
change := *req.GiftAmount - userInfo.GiftAmount
|
||||
if userInfo.GiftAmount != req.GiftAmount {
|
||||
change := req.GiftAmount - userInfo.GiftAmount
|
||||
if change != 0 {
|
||||
var changeType uint16
|
||||
if userInfo.GiftAmount < *req.GiftAmount {
|
||||
if userInfo.GiftAmount < req.GiftAmount {
|
||||
changeType = log.GiftTypeIncrease
|
||||
} else {
|
||||
changeType = log.GiftTypeReduce
|
||||
@@ -81,7 +81,7 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
||||
giftLog := log.Gift{
|
||||
Type: changeType,
|
||||
Amount: change,
|
||||
Balance: *req.GiftAmount,
|
||||
Balance: req.GiftAmount,
|
||||
Remark: "Admin adjustment",
|
||||
Timestamp: time.Now().UnixMilli(),
|
||||
}
|
||||
@@ -96,27 +96,23 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
userInfo.GiftAmount = *req.GiftAmount
|
||||
userInfo.GiftAmount = req.GiftAmount
|
||||
}
|
||||
}
|
||||
|
||||
if req.Commission != nil && *req.Commission != userInfo.Commission {
|
||||
remark := ""
|
||||
if req.Remark != nil {
|
||||
remark = *req.Remark
|
||||
}
|
||||
if isWithdrawalScene(remark) {
|
||||
if req.Commission != userInfo.Commission {
|
||||
if isWithdrawalScene(req.Remark) {
|
||||
logWithdrawalGuard(l.Logger, userInfo.Id)
|
||||
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 {
|
||||
return err
|
||||
}
|
||||
if err = logicCommon.WriteCommissionLog(tx, userInfo.Id, log.CommissionTypeAdjust, change, ""); err != nil {
|
||||
return err
|
||||
}
|
||||
userInfo.Commission = *req.Commission
|
||||
userInfo.Commission = req.Commission
|
||||
}
|
||||
if req.Avatar != "" {
|
||||
userInfo.Avatar = req.Avatar
|
||||
@@ -124,17 +120,15 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
||||
if req.ReferCode != "" {
|
||||
userInfo.ReferCode = req.ReferCode
|
||||
}
|
||||
if req.RefererId != nil {
|
||||
userInfo.RefererId = *req.RefererId
|
||||
}
|
||||
userInfo.RefererId = req.RefererId
|
||||
if req.Enable != nil {
|
||||
userInfo.Enable = req.Enable
|
||||
}
|
||||
if req.IsAdmin != nil {
|
||||
userInfo.IsAdmin = req.IsAdmin
|
||||
}
|
||||
if req.Remark != nil {
|
||||
userInfo.Remark = *req.Remark
|
||||
if req.Remark != "" {
|
||||
userInfo.Remark = req.Remark
|
||||
}
|
||||
if req.OnlyFirstPurchase != nil {
|
||||
userInfo.OnlyFirstPurchase = req.OnlyFirstPurchase
|
||||
|
||||
@@ -3226,17 +3226,17 @@ type UpdateUserBasiceInfoRequest struct {
|
||||
UserId int64 `json:"user_id" validate:"required"`
|
||||
Password string `json:"password"`
|
||||
Avatar string `json:"avatar"`
|
||||
Balance *int64 `json:"balance"`
|
||||
Commission *int64 `json:"commission"`
|
||||
Balance int64 `json:"balance"`
|
||||
Commission int64 `json:"commission"`
|
||||
ReferralPercentage uint8 `json:"referral_percentage"`
|
||||
OnlyFirstPurchase *bool `json:"only_first_purchase"`
|
||||
GiftAmount *int64 `json:"gift_amount"`
|
||||
GiftAmount int64 `json:"gift_amount"`
|
||||
Telegram int64 `json:"telegram"`
|
||||
ReferCode string `json:"refer_code"`
|
||||
RefererId *int64 `json:"referer_id"`
|
||||
RefererId int64 `json:"referer_id"`
|
||||
Enable *bool `json:"enable"`
|
||||
IsAdmin *bool `json:"is_admin"`
|
||||
Remark *string `json:"remark"`
|
||||
Remark string `json:"remark"`
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -48,4 +48,7 @@ func RegisterHandlers(mux *asynq.ServeMux, serverCtx *svc.ServiceContext) {
|
||||
// Apple IAP 对账(第二层:5min 扫描 + 第三层:日终全量)
|
||||
mux.Handle(types.SchedulerIAPReconcile, iapLogic.NewReconcileLogic(serverCtx))
|
||||
mux.Handle(types.SchedulerIAPDailyReconcile, iapLogic.NewDailyReconcileLogic(serverCtx))
|
||||
|
||||
// Stuck order recovery: reset claimed orders that timed out back to paid
|
||||
mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx))
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/server/internal/logic/admin/group"
|
||||
@@ -46,7 +47,7 @@ const (
|
||||
OrderStatusPaid = 2 // Order paid and ready for processing
|
||||
OrderStatusClose = 3 // Order closed/cancelled
|
||||
OrderStatusFailed = 4 // Order processing failed
|
||||
OrderStatusClaimed = 4 // Internal transient claim while a worker processes the order
|
||||
OrderStatusClaimed = 6 // Internal transient claim while a worker processes the order
|
||||
OrderStatusFinished = 5 // Order successfully completed
|
||||
)
|
||||
|
||||
@@ -94,6 +95,11 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
|
||||
if err != nil {
|
||||
// 如果订单不存在或状态不对,不重试
|
||||
if errors.Is(err, ErrInvalidOrderStatus) {
|
||||
if strings.Contains(err.Error(), "stuck in claimed") {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单卡在 claimed,将重试",
|
||||
logger.Field("order_no", payload.OrderNo))
|
||||
return err
|
||||
}
|
||||
logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过",
|
||||
logger.Field("order_no", payload.OrderNo))
|
||||
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 {
|
||||
l.releaseClaim(ctx, orderInfo.OrderNo)
|
||||
if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] releaseClaim also failed, stuck recovery will handle",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("release_error", releaseErr.Error()),
|
||||
)
|
||||
}
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] 处理订单失败,将重试",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("order_type", orderInfo.Type),
|
||||
@@ -125,7 +136,12 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
|
||||
}
|
||||
|
||||
if err = l.reconcilePostOrderSubscriptions(ctx, orderInfo); err != nil {
|
||||
l.releaseClaim(ctx, orderInfo.OrderNo)
|
||||
if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] releaseClaim also failed, stuck recovery will handle",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("release_error", releaseErr.Error()),
|
||||
)
|
||||
}
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单订阅兜底合并失败,将重试",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("order_type", orderInfo.Type),
|
||||
@@ -176,6 +192,14 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
|
||||
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 {
|
||||
logger.WithContext(ctx).Error("Order status error",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
@@ -201,7 +225,7 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
|
||||
return &orderInfo, nil
|
||||
}
|
||||
|
||||
func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) {
|
||||
func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) error {
|
||||
if err := l.svc.DB.WithContext(ctx).
|
||||
Model(&order.Order{}).
|
||||
Where("order_no = ? AND status = ?", orderNo, OrderStatusClaimed).
|
||||
@@ -210,7 +234,9 @@ func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) {
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("order_no", orderNo),
|
||||
)
|
||||
return fmt.Errorf("release claim failed for order %s: %w", orderNo, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// processOrderByType routes order processing based on the order type
|
||||
@@ -563,10 +589,26 @@ func (l *ActivateOrderLogic) finalizeCouponAndOrder(ctx context.Context, orderIn
|
||||
}
|
||||
}
|
||||
|
||||
// Update order status using state-guarded UpdateOrderStatus to prevent double finalization
|
||||
if err := l.svc.OrderModel.UpdateOrderStatus(ctx, orderInfo.OrderNo, OrderStatusFinished); err != nil {
|
||||
logger.WithContext(ctx).Error("Update order status failed",
|
||||
logger.Field("error", err.Error()),
|
||||
// Direct update from claimed(6) → finished(5), bypassing the model's status<target guard.
|
||||
// UpdateOrderStatus uses WHERE status < ?, so status=6 > 5 would never match.
|
||||
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 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),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -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"
|
||||
SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单
|
||||
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())
|
||||
}
|
||||
|
||||
// schedule stuck order recovery: every 10 minutes
|
||||
stuckOrderTask := asynq.NewTask(types.SchedulerStuckOrderRecovery, nil)
|
||||
if _, err := m.server.Register("@every 10m", stuckOrderTask, asynq.MaxRetry(1)); err != nil {
|
||||
logger.Errorf("register stuck order recovery task failed: %s", err.Error())
|
||||
}
|
||||
|
||||
if err := m.server.Run(); err != nil {
|
||||
logger.Errorf("run scheduler failed: %s", err.Error())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user