Merge remote-tracking branch 'origin/master' into internal
This commit is contained in:
@@ -0,0 +1,87 @@
|
||||
package group
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/server/internal/logic/admin/group"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
"github.com/perfect-panel/server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
)
|
||||
|
||||
type RecalculateGroupLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewRecalculateGroupLogic(svc *svc.ServiceContext) *RecalculateGroupLogic {
|
||||
return &RecalculateGroupLogic{
|
||||
svc: svc,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *RecalculateGroupLogic) ProcessTask(ctx context.Context, t *asynq.Task) error {
|
||||
logger.Infof("[RecalculateGroup] Starting scheduled group recalculation: %s", time.Now().Format("2006-01-02 15:04:05"))
|
||||
|
||||
// 1. Check if group management is enabled
|
||||
var enabledConfig struct {
|
||||
Value string `gorm:"column:value"`
|
||||
}
|
||||
err := l.svc.DB.Table("system").
|
||||
Where("`category` = ? AND `key` = ?", "group", "enabled").
|
||||
Select("value").
|
||||
First(&enabledConfig).Error
|
||||
if err != nil {
|
||||
logger.Errorw("[RecalculateGroup] Failed to read group enabled config", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
|
||||
// If not enabled, skip execution
|
||||
if enabledConfig.Value != "true" && enabledConfig.Value != "1" {
|
||||
logger.Debugf("[RecalculateGroup] Group management is not enabled, skipping")
|
||||
return nil
|
||||
}
|
||||
|
||||
// 2. Get grouping mode
|
||||
var modeConfig struct {
|
||||
Value string `gorm:"column:value"`
|
||||
}
|
||||
err = l.svc.DB.Table("system").
|
||||
Where("`category` = ? AND `key` = ?", "group", "mode").
|
||||
Select("value").
|
||||
First(&modeConfig).Error
|
||||
if err != nil {
|
||||
logger.Errorw("[RecalculateGroup] Failed to read group mode config", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
|
||||
mode := modeConfig.Value
|
||||
if mode == "" {
|
||||
mode = "average" // default mode
|
||||
}
|
||||
|
||||
// 3. Only execute if mode is "traffic"
|
||||
if mode != "traffic" {
|
||||
logger.Debugf("[RecalculateGroup] Group mode is not 'traffic' (current: %s), skipping", mode)
|
||||
return nil
|
||||
}
|
||||
|
||||
// 4. Execute traffic-based grouping
|
||||
logger.Infof("[RecalculateGroup] Executing traffic-based grouping")
|
||||
|
||||
logic := group.NewRecalculateGroupLogic(ctx, l.svc)
|
||||
req := &types.RecalculateGroupRequest{
|
||||
Mode: "traffic",
|
||||
TriggerType: "scheduled",
|
||||
}
|
||||
|
||||
if err := logic.RecalculateGroup(req); err != nil {
|
||||
logger.Errorw("[RecalculateGroup] Failed to execute traffic grouping", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
|
||||
logger.Infof("[RecalculateGroup] Successfully completed traffic-based grouping: %s", time.Now().Format("2006-01-02 15:04:05"))
|
||||
return nil
|
||||
}
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/server/internal/logic/admin/group"
|
||||
"github.com/perfect-panel/server/internal/model/log"
|
||||
"github.com/perfect-panel/server/pkg/constant"
|
||||
"github.com/perfect-panel/server/pkg/logger"
|
||||
@@ -24,9 +25,10 @@ import (
|
||||
"github.com/perfect-panel/server/internal/model/subscribe"
|
||||
"github.com/perfect-panel/server/internal/model/user"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
"github.com/perfect-panel/server/pkg/tool"
|
||||
"github.com/perfect-panel/server/pkg/uuidx"
|
||||
"github.com/perfect-panel/server/queue/types"
|
||||
queueTypes "github.com/perfect-panel/server/queue/types"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
@@ -126,8 +128,8 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
|
||||
}
|
||||
|
||||
// parsePayload unMarshals the task payload into a structured format
|
||||
func (l *ActivateOrderLogic) parsePayload(ctx context.Context, payload []byte) (*types.ForthwithActivateOrderPayload, error) {
|
||||
var p types.ForthwithActivateOrderPayload
|
||||
func (l *ActivateOrderLogic) parsePayload(ctx context.Context, payload []byte) (*queueTypes.ForthwithActivateOrderPayload, error) {
|
||||
var p queueTypes.ForthwithActivateOrderPayload
|
||||
if err := json.Unmarshal(payload, &p); err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Unmarshal payload failed",
|
||||
logger.Field("error", err.Error()),
|
||||
@@ -322,6 +324,9 @@ func (l *ActivateOrderLogic) NewPurchase(ctx context.Context, orderInfo *order.O
|
||||
}
|
||||
}
|
||||
|
||||
// Trigger user group recalculation (runs in background)
|
||||
l.triggerUserGroupRecalculation(ctx, userInfo.Id)
|
||||
|
||||
// Handle commission in separate goroutine to avoid blocking
|
||||
go l.handleCommission(context.Background(), userInfo, orderInfo)
|
||||
|
||||
@@ -782,6 +787,63 @@ func (l *ActivateOrderLogic) clearServerCache(ctx context.Context, sub *subscrib
|
||||
}
|
||||
}
|
||||
|
||||
// triggerUserGroupRecalculation triggers user group recalculation after subscription changes
|
||||
// This runs asynchronously in background to avoid blocking the main order processing flow
|
||||
func (l *ActivateOrderLogic) triggerUserGroupRecalculation(ctx context.Context, userId int64) {
|
||||
go func() {
|
||||
// Use a new context with timeout for group recalculation
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// Check if group management is enabled
|
||||
var groupEnabled string
|
||||
err := l.svc.DB.Table("system").
|
||||
Where("`category` = ? AND `key` = ?", "group", "enabled").
|
||||
Select("value").
|
||||
Scan(&groupEnabled).Error
|
||||
if err != nil || groupEnabled != "true" && groupEnabled != "1" {
|
||||
logger.Debugf("[Group Trigger] Group management not enabled, skipping recalculation")
|
||||
return
|
||||
}
|
||||
|
||||
// Get the configured grouping mode
|
||||
var groupMode string
|
||||
err = l.svc.DB.Table("system").
|
||||
Where("`category` = ? AND `key` = ?", "group", "mode").
|
||||
Select("value").
|
||||
Scan(&groupMode).Error
|
||||
if err != nil {
|
||||
logger.Errorw("[Group Trigger] Failed to get group mode", logger.Field("error", err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
// Validate group mode
|
||||
if groupMode != "average" && groupMode != "subscribe" && groupMode != "traffic" {
|
||||
logger.Debugf("[Group Trigger] Invalid group mode (current: %s), skipping", groupMode)
|
||||
return
|
||||
}
|
||||
|
||||
// Trigger group recalculation with the configured mode
|
||||
logic := group.NewRecalculateGroupLogic(ctx, l.svc)
|
||||
req := &types.RecalculateGroupRequest{
|
||||
Mode: groupMode,
|
||||
}
|
||||
|
||||
if err := logic.RecalculateGroup(req); err != nil {
|
||||
logger.Errorw("[Group Trigger] Failed to recalculate user group",
|
||||
logger.Field("user_id", userId),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return
|
||||
}
|
||||
|
||||
logger.Infow("[Group Trigger] Successfully recalculated user group",
|
||||
logger.Field("user_id", userId),
|
||||
logger.Field("mode", groupMode),
|
||||
)
|
||||
}()
|
||||
}
|
||||
|
||||
// Renewal handles subscription renewal including subscription extension,
|
||||
// traffic reset (if configured), commission processing, and notifications
|
||||
func (l *ActivateOrderLogic) Renewal(ctx context.Context, orderInfo *order.Order, iapExpireAt int64) error {
|
||||
@@ -907,6 +969,9 @@ func (l *ActivateOrderLogic) updateSubscriptionForRenewal(ctx context.Context, u
|
||||
|
||||
userSub.ExpireTime = tool.AddTime(sub.UnitTime, orderInfo.Quantity, userSub.ExpireTime)
|
||||
userSub.Status = 1
|
||||
// 续费时重置过期流量字段
|
||||
userSub.ExpiredDownload = 0
|
||||
userSub.ExpiredUpload = 0
|
||||
|
||||
if err := l.svc.UserModel.UpdateSubscribe(ctx, userSub); err != nil {
|
||||
logger.WithContext(ctx).Error("Update user subscribe failed", logger.Field("error", err.Error()))
|
||||
@@ -931,6 +996,8 @@ func (l *ActivateOrderLogic) ResetTraffic(ctx context.Context, orderInfo *order.
|
||||
// Reset traffic
|
||||
userSub.Download = 0
|
||||
userSub.Upload = 0
|
||||
userSub.ExpiredDownload = 0
|
||||
userSub.ExpiredUpload = 0
|
||||
userSub.Status = 1
|
||||
|
||||
if err := l.svc.UserModel.UpdateSubscribe(ctx, userSub); err != nil {
|
||||
@@ -1242,6 +1309,7 @@ func (l *ActivateOrderLogic) RedemptionActivate(ctx context.Context, orderInfo *
|
||||
Traffic: us.Traffic,
|
||||
Download: us.Download,
|
||||
Upload: us.Upload,
|
||||
NodeGroupId: us.NodeGroupId,
|
||||
}
|
||||
break
|
||||
}
|
||||
@@ -1328,6 +1396,7 @@ func (l *ActivateOrderLogic) RedemptionActivate(ctx context.Context, orderInfo *
|
||||
Token: uuidx.SubscribeToken(orderInfo.OrderNo),
|
||||
UUID: uuid.New().String(),
|
||||
Status: 1,
|
||||
NodeGroupId: sub.NodeGroupId, // Inherit node_group_id from subscription plan
|
||||
}
|
||||
|
||||
err = l.svc.UserModel.InsertSubscribe(ctx, newSubscribe, tx)
|
||||
@@ -1374,6 +1443,9 @@ func (l *ActivateOrderLogic) RedemptionActivate(ctx context.Context, orderInfo *
|
||||
return err
|
||||
}
|
||||
|
||||
// Trigger user group recalculation (runs in background)
|
||||
l.triggerUserGroupRecalculation(ctx, userInfo.Id)
|
||||
|
||||
// 7. 清理缓存(关键步骤:让节点获取最新订阅)
|
||||
l.clearServerCache(ctx, sub)
|
||||
|
||||
|
||||
@@ -62,7 +62,6 @@ func (l *CheckSubscriptionLogic) ProcessTask(ctx context.Context, _ *asynq.Task)
|
||||
}
|
||||
l.clearServerCache(ctx, list...)
|
||||
logger.Infow("[Check Subscription Traffic] Update subscribe status", logger.Field("user_ids", ids), logger.Field("count", int64(len(ids))))
|
||||
|
||||
} else {
|
||||
logger.Info("[Check Subscription Traffic] No subscribe need to update")
|
||||
}
|
||||
@@ -108,6 +107,7 @@ func (l *CheckSubscriptionLogic) ProcessTask(ctx context.Context, _ *asynq.Task)
|
||||
} else {
|
||||
logger.Info("[Check Subscription Expire] No subscribe need to update")
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
|
||||
@@ -98,11 +98,13 @@ func (l *TrafficStatisticsLogic) ProcessTask(ctx context.Context, task *asynq.Ta
|
||||
// update user subscribe with log
|
||||
d := int64(float32(log.Download) * ratio * realTimeMultiplier)
|
||||
u := int64(float32(log.Upload) * ratio * realTimeMultiplier)
|
||||
if err := l.svc.UserModel.UpdateUserSubscribeWithTraffic(ctx, sub.Id, d, u); err != nil {
|
||||
isExpired := now.After(sub.ExpireTime)
|
||||
if err := l.svc.UserModel.UpdateUserSubscribeWithTraffic(ctx, sub.Id, d, u, isExpired); err != nil {
|
||||
logger.WithContext(ctx).Error("[TrafficStatistics] Update user subscribe with log failed",
|
||||
logger.Field("sid", log.SID),
|
||||
logger.Field("download", float32(log.Download)*ratio),
|
||||
logger.Field("upload", float32(log.Upload)*ratio),
|
||||
logger.Field("is_expired", isExpired),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
continue
|
||||
|
||||
Reference in New Issue
Block a user