feat(api): migrate server and node data handling, update related structures and logic

This commit is contained in:
Chang lue Tsen
2025-08-26 07:05:59 -04:00
parent 9b3cdbbb4f
commit c7884d94aa
52 changed files with 1079 additions and 458 deletions
-38
View File
@@ -2,14 +2,9 @@ package countrylogic
import (
"context"
"encoding/json"
"github.com/perfect-panel/server/pkg/logger"
"github.com/hibiken/asynq"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/pkg/ip"
"github.com/perfect-panel/server/queue/types"
)
type GetNodeCountryLogic struct {
@@ -22,39 +17,6 @@ func NewGetNodeCountryLogic(svcCtx *svc.ServiceContext) *GetNodeCountryLogic {
}
}
func (l *GetNodeCountryLogic) ProcessTask(ctx context.Context, task *asynq.Task) error {
var payload types.GetNodeCountry
if err := json.Unmarshal(task.Payload(), &payload); err != nil {
logger.WithContext(ctx).Error("[GetNodeCountryLogic] Unmarshal payload failed",
logger.Field("error", err.Error()),
logger.Field("payload", task.Payload()),
)
return nil
}
serverAddr := payload.ServerAddr
resp, err := ip.GetRegionByIp(serverAddr)
if err != nil {
logger.WithContext(ctx).Error("[GetNodeCountryLogic] ", logger.Field("error", err.Error()), logger.Field("serverAddr", serverAddr))
return nil
}
servers, err := l.svcCtx.ServerModel.FindNodeByServerAddrAndProtocol(ctx, payload.ServerAddr, payload.Protocol)
if err != nil {
logger.WithContext(ctx).Error("[GetNodeCountryLogic] FindNodeByServerAddrAnd", logger.Field("error", err.Error()), logger.Field("serverAddr", serverAddr))
return err
}
if len(servers) == 0 {
return nil
}
for _, ser := range servers {
ser.Country = resp.Country
ser.City = resp.City
ser.Latitude = resp.Latitude
ser.Longitude = resp.Longitude
err := l.svcCtx.ServerModel.Update(ctx, ser)
if err != nil {
logger.WithContext(ctx).Error("[GetNodeCountryLogic] ", logger.Field("error", err.Error()), logger.Field("id", ser.Id))
}
}
logger.WithContext(ctx).Info("[GetNodeCountryLogic] ", logger.Field("country", resp.Country), logger.Field("city", resp.Country))
return nil
}
+11 -26
View File
@@ -7,16 +7,17 @@ import (
"encoding/json"
"fmt"
"strconv"
"strings"
"time"
"github.com/perfect-panel/server/internal/model/log"
"github.com/perfect-panel/server/internal/model/node"
"github.com/perfect-panel/server/pkg/constant"
"github.com/perfect-panel/server/pkg/logger"
tgbotapi "github.com/go-telegram-bot-api/telegram-bot-api/v5"
"github.com/google/uuid"
"github.com/hibiken/asynq"
"github.com/perfect-panel/server/internal/config"
"github.com/perfect-panel/server/internal/logic/telegram"
"github.com/perfect-panel/server/internal/model/order"
"github.com/perfect-panel/server/internal/model/subscribe"
@@ -429,34 +430,18 @@ func (l *ActivateOrderLogic) calculateCommission(price int64) int64 {
// clearServerCache clears user list cache for all servers associated with the subscription
func (l *ActivateOrderLogic) clearServerCache(ctx context.Context, sub *subscribe.Subscribe) {
serverIds := tool.StringToInt64Slice(sub.Server)
groupServerIds := l.getServerIdsByGroups(ctx, sub.ServerGroup)
allServerIds := append(serverIds, groupServerIds...)
nodeIds := tool.StringToInt64Slice(sub.Nodes)
tags := strings.Split(sub.NodeTags, ",")
for _, id := range allServerIds {
cacheKey := fmt.Sprintf("%s%d", config.ServerUserListCacheKey, id)
if err := l.svc.Redis.Del(ctx, cacheKey).Err(); err != nil {
logger.WithContext(ctx).Error("Del server user list cache failed",
logger.Field("error", err.Error()),
logger.Field("cache_key", cacheKey),
)
}
}
}
// getServerIdsByGroups retrieves server IDs from server groups
func (l *ActivateOrderLogic) getServerIdsByGroups(ctx context.Context, serverGroup string) []int64 {
data, err := l.svc.ServerModel.FindServerListByGroupIds(ctx, tool.StringToInt64Slice(serverGroup))
err := l.svc.NodeModel.ClearNodeCache(ctx, &node.FilterNodeParams{
Page: 1,
Size: 1000,
ServerId: nodeIds,
Tag: tags,
})
if err != nil {
logger.WithContext(ctx).Error("Find server list failed", logger.Field("error", err.Error()))
return nil
logger.WithContext(ctx).Error("[Order Queue] Clear node cache failed", logger.Field("error", err.Error()))
}
serverIds := make([]int64, len(data))
for i, item := range data {
serverIds[i] = item.Id
}
return serverIds
}
// Renewal handles subscription renewal including subscription extension,
+2 -2
View File
@@ -73,7 +73,7 @@ func (l *ServerDataLogic) getRanking(ctx context.Context) (top10ServerToday, top
if s.ServerId == 0 {
continue
}
serverInfo, err := l.svc.ServerModel.FindOne(ctx, s.ServerId)
serverInfo, err := l.svc.NodeModel.FindOneServer(ctx, s.ServerId)
if err != nil {
logger.Error("[ServerDataLogic] Find server failed", logger.Field("error", err.Error()))
continue
@@ -92,7 +92,7 @@ func (l *ServerDataLogic) getRanking(ctx context.Context) (top10ServerToday, top
logger.Error("[ServerDataLogic] Get top servers traffic by day failed", logger.Field("error", err.Error()))
} else {
for _, s := range serverYesterday {
serverInfo, err := l.svc.ServerModel.FindOne(ctx, s.ServerId)
serverInfo, err := l.svc.NodeModel.FindOneServer(ctx, s.ServerId)
if err != nil {
logger.Error("[ServerDataLogic] Find server failed", logger.Field("error", err.Error()))
continue
+11 -12
View File
@@ -38,7 +38,7 @@ func (l *TrafficStatisticsLogic) ProcessTask(ctx context.Context, task *asynq.Ta
return nil
}
// query server info
serverInfo, err := l.svc.ServerModel.FindOne(ctx, payload.ServerId)
serverInfo, err := l.svc.NodeModel.FindOneServer(ctx, payload.ServerId)
if err != nil {
logger.WithContext(ctx).Error("[TrafficStatistics] Find server info failed",
logger.Field("serverId", payload.ServerId),
@@ -46,23 +46,22 @@ func (l *TrafficStatisticsLogic) ProcessTask(ctx context.Context, task *asynq.Ta
)
return nil
}
if serverInfo.TrafficRatio == 0 {
logger.WithContext(ctx).Error("[TrafficStatistics] Server log ratio is 0",
logger.Field("serverId", payload.ServerId),
)
return nil
var serverRatio float32 = 1.0
if serverInfo.Ratio > 0 {
serverRatio = serverInfo.Ratio
}
now := time.Now()
realTimeMultiplier := l.svc.NodeMultiplierManager.GetMultiplier(now)
for _, log := range payload.Logs {
// update user subscribe with log
d := int64(float32(log.Download) * serverInfo.TrafficRatio * realTimeMultiplier)
u := int64(float32(log.Upload) * serverInfo.TrafficRatio * realTimeMultiplier)
d := int64(float32(log.Download) * serverRatio * realTimeMultiplier)
u := int64(float32(log.Upload) * serverRatio * realTimeMultiplier)
if err := l.svc.UserModel.UpdateUserSubscribeWithTraffic(ctx, log.SID, d, u); err != nil {
logger.WithContext(ctx).Error("[TrafficStatistics] Update user subscribe with log failed",
logger.Field("sid", log.SID),
logger.Field("download", float32(log.Download)*serverInfo.TrafficRatio),
logger.Field("upload", float32(log.Upload)*serverInfo.TrafficRatio),
logger.Field("download", float32(log.Download)*serverRatio),
logger.Field("upload", float32(log.Upload)*serverRatio),
logger.Field("error", err.Error()),
)
continue
@@ -88,8 +87,8 @@ func (l *TrafficStatisticsLogic) ProcessTask(ctx context.Context, task *asynq.Ta
}); err != nil {
logger.WithContext(ctx).Error("[TrafficStatistics] Create log log failed",
logger.Field("uid", log.SID),
logger.Field("download", float32(log.Download)*serverInfo.TrafficRatio),
logger.Field("upload", float32(log.Upload)*serverInfo.TrafficRatio),
logger.Field("download", float32(log.Download)*serverRatio),
logger.Field("upload", float32(log.Upload)*serverRatio),
logger.Field("error", err.Error()),
)
}