init
This commit is contained in:
@@ -0,0 +1,60 @@
|
||||
package countrylogic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/ip"
|
||||
"github.com/perfect-panel/ppanel-server/queue/types"
|
||||
)
|
||||
|
||||
type GetNodeCountryLogic struct {
|
||||
svcCtx *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewGetNodeCountryLogic(svcCtx *svc.ServiceContext) *GetNodeCountryLogic {
|
||||
return &GetNodeCountryLogic{
|
||||
svcCtx: svcCtx,
|
||||
}
|
||||
}
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
package emailLogic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/perfect-panel/ppanel-server/internal/model/log"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/email"
|
||||
"github.com/perfect-panel/ppanel-server/queue/types"
|
||||
)
|
||||
|
||||
type SendEmailLogic struct {
|
||||
svcCtx *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewSendEmailLogic(svcCtx *svc.ServiceContext) *SendEmailLogic {
|
||||
return &SendEmailLogic{
|
||||
svcCtx: svcCtx,
|
||||
}
|
||||
}
|
||||
func (l *SendEmailLogic) ProcessTask(ctx context.Context, task *asynq.Task) error {
|
||||
var payload types.SendEmailPayload
|
||||
if err := json.Unmarshal(task.Payload(), &payload); err != nil {
|
||||
logger.WithContext(ctx).Error("[SendEmailLogic] Unmarshal payload failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("payload", task.Payload()),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
messageLog := log.MessageLog{
|
||||
Type: log.Email.String(),
|
||||
Platform: l.svcCtx.Config.Email.Platform,
|
||||
To: payload.Email,
|
||||
Subject: payload.Subject,
|
||||
Content: payload.Content,
|
||||
}
|
||||
sender, err := email.NewSender(l.svcCtx.Config.Email.Platform, l.svcCtx.Config.Email.PlatformConfig, l.svcCtx.Config.Site.SiteName)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[SendEmailLogic] NewSender failed", logger.Field("error", err.Error()))
|
||||
return nil
|
||||
}
|
||||
err = sender.Send([]string{payload.Email}, payload.Subject, payload.Content)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[SendEmailLogic] Send email failed", logger.Field("error", err.Error()))
|
||||
return nil
|
||||
}
|
||||
messageLog.Status = 1
|
||||
if err = l.svcCtx.LogModel.InsertMessageLog(ctx, &messageLog); err != nil {
|
||||
logger.WithContext(ctx).Error("[SendEmailLogic] InsertMessageLog failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("messageLog", messageLog),
|
||||
)
|
||||
}
|
||||
logger.WithContext(ctx).Info("[SendEmailLogic] Send email", logger.Field("email", payload.Email), logger.Field("content", payload.Content))
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,691 @@
|
||||
package orderLogic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/constant"
|
||||
|
||||
"github.com/perfect-panel/ppanel-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/ppanel-server/internal/config"
|
||||
"github.com/perfect-panel/ppanel-server/internal/logic/telegram"
|
||||
"github.com/perfect-panel/ppanel-server/internal/model/order"
|
||||
"github.com/perfect-panel/ppanel-server/internal/model/user"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/tool"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/uuidx"
|
||||
"github.com/perfect-panel/ppanel-server/queue/types"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const (
|
||||
Subscribe = 1
|
||||
Renewal = 2
|
||||
ResetTraffic = 3
|
||||
Recharge = 4
|
||||
)
|
||||
|
||||
type ActivateOrderLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewActivateOrderLogic(svc *svc.ServiceContext) *ActivateOrderLogic {
|
||||
return &ActivateOrderLogic{
|
||||
svc: svc,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task) error {
|
||||
payload := types.ForthwithActivateOrderPayload{}
|
||||
if err := json.Unmarshal(task.Payload(), &payload); err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Unmarshal payload failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("payload", string(task.Payload())),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
// Find order by order no
|
||||
orderInfo, err := l.svc.OrderModel.FindOneByOrderNo(ctx, payload.OrderNo)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find order failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("order_no", payload.OrderNo),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
|
||||
if orderInfo.Status != 2 {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Order status error",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("status", orderInfo.Status),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
switch orderInfo.Type {
|
||||
case Subscribe:
|
||||
err = l.NewPurchase(ctx, orderInfo)
|
||||
case Renewal:
|
||||
err = l.Renewal(ctx, orderInfo)
|
||||
case ResetTraffic:
|
||||
err = l.ResetTraffic(ctx, orderInfo)
|
||||
case Recharge:
|
||||
err = l.Recharge(ctx, orderInfo)
|
||||
default:
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Order type is invalid", logger.Field("type", orderInfo.Type))
|
||||
}
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Process task failed", logger.Field("error", err.Error()))
|
||||
return nil
|
||||
}
|
||||
// if coupon is not empty
|
||||
if orderInfo.Coupon != "" {
|
||||
// update coupon status
|
||||
err = l.svc.CouponModel.UpdateCount(ctx, orderInfo.Coupon)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Update coupon status failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("coupon", orderInfo.Coupon),
|
||||
)
|
||||
}
|
||||
}
|
||||
// update order status
|
||||
orderInfo.Status = 5
|
||||
err = l.svc.OrderModel.Update(ctx, orderInfo)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Update order status failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// NewPurchase New purchase
|
||||
func (l *ActivateOrderLogic) NewPurchase(ctx context.Context, orderInfo *order.Order) error {
|
||||
var userInfo *user.User
|
||||
var err error
|
||||
if orderInfo.UserId != 0 {
|
||||
// find user by user id
|
||||
userInfo, err = l.svc.UserModel.FindOne(ctx, orderInfo.UserId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find user failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("user_id", orderInfo.UserId),
|
||||
logger.Field("user_id", orderInfo.UserId),
|
||||
)
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
// If User ID is 0, it means that the order is a guest order, need to create a new user
|
||||
// query info with redis
|
||||
cacheKey := fmt.Sprintf(constant.TempOrderCacheKey, orderInfo.OrderNo)
|
||||
data, err := l.svc.Redis.Get(ctx, cacheKey).Result()
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Get temp order cache failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("cache_key", cacheKey),
|
||||
)
|
||||
return err
|
||||
}
|
||||
var tempOrder constant.TemporaryOrderInfo
|
||||
if err = json.Unmarshal([]byte(data), &tempOrder); err != nil {
|
||||
logger.WithContext(ctx).Errorw("[ActivateOrderLogic] Unmarshal temp order failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// create user
|
||||
|
||||
userInfo = &user.User{
|
||||
Password: tool.EncodePassWord(tempOrder.Password),
|
||||
AuthMethods: []user.AuthMethods{
|
||||
{
|
||||
AuthType: tempOrder.AuthType,
|
||||
AuthIdentifier: tempOrder.Identifier,
|
||||
},
|
||||
},
|
||||
}
|
||||
err = l.svc.UserModel.Transaction(ctx, func(tx *gorm.DB) error {
|
||||
// Save user information
|
||||
if err := tx.Save(userInfo).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
// Generate ReferCode
|
||||
userInfo.ReferCode = uuidx.UserInviteCode(userInfo.Id)
|
||||
// Update ReferCode
|
||||
if err := tx.Model(&user.User{}).Where("id = ?", userInfo.Id).Update("refer_code", userInfo.ReferCode).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
orderInfo.UserId = userInfo.Id
|
||||
return tx.Model(&order.Order{}).Where("order_no = ?", orderInfo.OrderNo).Update("user_id", userInfo.Id).Error
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Create user failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
logger.WithContext(ctx).Info("[ActivateOrderLogic] Create guest user success", logger.Field("user_id", userInfo.Id), logger.Field("Identifier", tempOrder.Identifier), logger.Field("AuthType", tempOrder.AuthType))
|
||||
}
|
||||
// find subscribe by id
|
||||
sub, err := l.svc.SubscribeModel.FindOne(ctx, orderInfo.SubscribeId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Errorw("[ActivateOrderLogic] Find subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("subscribe_id", orderInfo.SubscribeId),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// create user subscribe
|
||||
now := time.Now()
|
||||
|
||||
//系统开启了试用订阅功能,并且当前存在试用订阅
|
||||
if l.svc.Config.Register.EnableTrial && l.svc.Config.Register.TrialSubscribe != 0 {
|
||||
//查询使用订阅套餐
|
||||
subscribeDetails, subErr := l.svc.UserModel.QueryUserSubscribe(ctx, userInfo.Id, 1, 2)
|
||||
if subErr != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] disable user try out subscribe failed",
|
||||
logger.Field("error", subErr.Error()),
|
||||
)
|
||||
} else {
|
||||
for _, item := range subscribeDetails {
|
||||
if item.Subscribe.Id == l.svc.Config.Register.TrialSubscribe {
|
||||
err = l.svc.UserModel.UpdateSubscribe(ctx, &user.Subscribe{
|
||||
Id: item.Id,
|
||||
UserId: item.UserId,
|
||||
OrderId: item.OrderId,
|
||||
SubscribeId: item.SubscribeId,
|
||||
StartTime: item.StartTime,
|
||||
ExpireTime: now,
|
||||
Traffic: item.Traffic,
|
||||
Download: item.Download,
|
||||
Upload: item.Upload,
|
||||
Token: item.Token,
|
||||
UUID: item.UUID,
|
||||
FinishedAt: now,
|
||||
Status: 3,
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] disable user try out subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
} else {
|
||||
logger.WithContext(ctx).Info("[ActivateOrderLogic] disable user try out subscribe success",
|
||||
logger.Field("user_id", userInfo.Id),
|
||||
logger.Field("subscribe_id", item.SubscribeId),
|
||||
)
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
userSub := user.Subscribe{
|
||||
Id: 0,
|
||||
UserId: orderInfo.UserId,
|
||||
OrderId: orderInfo.Id,
|
||||
SubscribeId: orderInfo.SubscribeId,
|
||||
StartTime: now,
|
||||
ExpireTime: tool.AddTime(sub.UnitTime, orderInfo.Quantity, now),
|
||||
Traffic: sub.Traffic,
|
||||
Download: 0,
|
||||
Upload: 0,
|
||||
Token: uuidx.SubscribeToken(orderInfo.OrderNo),
|
||||
UUID: uuid.New().String(),
|
||||
Status: 1,
|
||||
}
|
||||
err = l.svc.UserModel.InsertSubscribe(ctx, &userSub)
|
||||
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Insert user subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// handler commission
|
||||
if userInfo.RefererId != 0 &&
|
||||
l.svc.Config.Invite.ReferralPercentage != 0 &&
|
||||
(!l.svc.Config.Invite.OnlyFirstPurchase || orderInfo.IsNew) {
|
||||
referer, err := l.svc.UserModel.FindOne(ctx, userInfo.RefererId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find referer failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("referer_id", userInfo.RefererId),
|
||||
)
|
||||
goto updateCache
|
||||
}
|
||||
// calculate commission
|
||||
amount := float64(orderInfo.Price) * (float64(l.svc.Config.Invite.ReferralPercentage) / 100)
|
||||
referer.Commission += int64(amount)
|
||||
err = l.svc.UserModel.Update(ctx, referer)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Update referer commission failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
goto updateCache
|
||||
}
|
||||
// create commission log
|
||||
commissionLog := user.CommissionLog{
|
||||
UserId: referer.Id,
|
||||
OrderNo: orderInfo.OrderNo,
|
||||
Amount: int64(amount),
|
||||
}
|
||||
err = l.svc.UserModel.InsertCommissionLog(ctx, &commissionLog)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Insert commission log failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
err = l.svc.UserModel.UpdateUserCache(ctx, referer)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Errorw("[ActivateOrderLogic] Update referer cache", logger.Field("error", err.Error()), logger.Field("user_id", referer.Id))
|
||||
}
|
||||
}
|
||||
updateCache:
|
||||
for _, id := range tool.StringToInt64Slice(sub.Server) {
|
||||
cacheKey := fmt.Sprintf("%s%d", config.ServerUserListCacheKey, id)
|
||||
err = l.svc.Redis.Del(ctx, cacheKey).Err()
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Del server user list cache failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("cache_key", cacheKey),
|
||||
)
|
||||
}
|
||||
}
|
||||
data, err := l.svc.ServerModel.FindServerListByGroupIds(ctx, tool.StringToInt64Slice(sub.ServerGroup))
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find server list failed", logger.Field("error", err.Error()))
|
||||
return nil
|
||||
}
|
||||
for _, item := range data {
|
||||
cacheKey := fmt.Sprintf("%s%d", config.ServerUserListCacheKey, item.Id)
|
||||
err = l.svc.Redis.Del(ctx, cacheKey).Err()
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Del server user list cache failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("cache_key", cacheKey),
|
||||
)
|
||||
}
|
||||
}
|
||||
userTelegramChatId, ok := findTelegram(userInfo)
|
||||
|
||||
// sendMessage To Telegram
|
||||
if ok {
|
||||
text, err := tool.RenderTemplateToString(telegram.PurchaseNotify, map[string]string{
|
||||
"OrderNo": orderInfo.OrderNo,
|
||||
"SubscribeName": sub.Name,
|
||||
"OrderAmount": fmt.Sprintf("%.2f", float64(orderInfo.Price)/100),
|
||||
"ExpireTime": userSub.ExpireTime.Format("2006-01-02 15:04:05"),
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendUserNotifyWithTelegram(userTelegramChatId, text)
|
||||
}
|
||||
// send message to admin
|
||||
text, err := tool.RenderTemplateToString(telegram.AdminOrderNotify, map[string]string{
|
||||
"OrderNo": orderInfo.OrderNo,
|
||||
"TradeNo": orderInfo.TradeNo,
|
||||
"SubscribeName": sub.Name,
|
||||
//"UserEmail": userInfo.Email,
|
||||
"OrderAmount": fmt.Sprintf("%.2f", float64(orderInfo.Price)/100),
|
||||
"OrderStatus": "已支付",
|
||||
"OrderTime": orderInfo.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||
"PaymentMethod": orderInfo.Method,
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render AdminOrderNotify template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendAdminNotifyWithTelegram(ctx, text)
|
||||
logger.WithContext(ctx).Info("[ActivateOrderLogic] Insert user subscribe success")
|
||||
return nil
|
||||
}
|
||||
|
||||
// Renewal Renewal
|
||||
func (l *ActivateOrderLogic) Renewal(ctx context.Context, orderInfo *order.Order) error {
|
||||
// find user by user id
|
||||
userInfo, err := l.svc.UserModel.FindOne(ctx, orderInfo.UserId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find user failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("user_id", orderInfo.UserId),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// find user subscribe by subscribe token
|
||||
userSub, err := l.svc.UserModel.FindOneSubscribeByOrderId(ctx, orderInfo.ParentId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find user subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("order_id", orderInfo.Id),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// find subscribe by id
|
||||
sub, err := l.svc.SubscribeModel.FindOne(ctx, orderInfo.SubscribeId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("subscribe_id", orderInfo.SubscribeId),
|
||||
logger.Field("order_id", orderInfo.Id),
|
||||
)
|
||||
return err
|
||||
}
|
||||
now := time.Now()
|
||||
if userSub.ExpireTime.Before(now) {
|
||||
userSub.ExpireTime = now
|
||||
userSub.Status = 1
|
||||
}
|
||||
|
||||
//fix bug:FinishedAt causes the update subscription to fail
|
||||
if now.AddDate(-30, 0, 0).After(userSub.FinishedAt) {
|
||||
userSub.FinishedAt = now
|
||||
}
|
||||
|
||||
userSub.ExpireTime = tool.AddTime(sub.UnitTime, orderInfo.Quantity, userSub.ExpireTime)
|
||||
// update user subscribe
|
||||
err = l.svc.UserModel.UpdateSubscribe(ctx, userSub)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Update user subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// handler commission
|
||||
if userInfo.RefererId != 0 &&
|
||||
l.svc.Config.Invite.ReferralPercentage != 0 &&
|
||||
!l.svc.Config.Invite.OnlyFirstPurchase {
|
||||
referer, err := l.svc.UserModel.FindOne(ctx, userInfo.RefererId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find referer failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("referer_id", userInfo.RefererId),
|
||||
)
|
||||
goto sendMessage
|
||||
}
|
||||
// calculate commission
|
||||
amount := float64(orderInfo.Price) * (float64(l.svc.Config.Invite.ReferralPercentage) / 100)
|
||||
referer.Commission += int64(amount)
|
||||
err = l.svc.UserModel.Update(ctx, referer)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Update referer commission failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
goto sendMessage
|
||||
}
|
||||
// create commission log
|
||||
commissionLog := user.CommissionLog{
|
||||
UserId: referer.Id,
|
||||
OrderNo: orderInfo.OrderNo,
|
||||
Amount: int64(amount),
|
||||
}
|
||||
err = l.svc.UserModel.InsertCommissionLog(ctx, &commissionLog)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Insert commission log failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
err = l.svc.UserModel.UpdateUserCache(ctx, referer)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Errorw("[ActivateOrderLogic] Update referer cache", logger.Field("error", err.Error()), logger.Field("user_id", referer.Id))
|
||||
}
|
||||
}
|
||||
sendMessage:
|
||||
userTelegramChatId, ok := findTelegram(userInfo)
|
||||
// SendMessage To Telegram
|
||||
if ok {
|
||||
text, err := tool.RenderTemplateToString(telegram.RenewalNotify, map[string]string{
|
||||
"OrderNo": orderInfo.OrderNo,
|
||||
"SubscribeName": sub.Name,
|
||||
"OrderAmount": fmt.Sprintf("%.2f", float64(orderInfo.Price)/100),
|
||||
"ExpireTime": userSub.ExpireTime.Format("2006-01-02 15:04:05"),
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendUserNotifyWithTelegram(userTelegramChatId, text)
|
||||
}
|
||||
|
||||
// send message to admin
|
||||
text, err := tool.RenderTemplateToString(telegram.AdminOrderNotify, map[string]string{
|
||||
"OrderNo": orderInfo.OrderNo,
|
||||
"TradeNo": orderInfo.TradeNo,
|
||||
"SubscribeName": sub.Name,
|
||||
//"UserEmail": userInfo.Email,
|
||||
"OrderAmount": fmt.Sprintf("%.2f", float64(orderInfo.Price)/100),
|
||||
"OrderStatus": "已支付",
|
||||
"OrderTime": orderInfo.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||
"PaymentMethod": orderInfo.Method,
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render AdminOrderNotify template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendAdminNotifyWithTelegram(ctx, text)
|
||||
return nil
|
||||
}
|
||||
|
||||
// ResetTraffic Reset traffic
|
||||
func (l *ActivateOrderLogic) ResetTraffic(ctx context.Context, orderInfo *order.Order) error {
|
||||
// find user by user id
|
||||
userInfo, err := l.svc.UserModel.FindOne(ctx, orderInfo.UserId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find user failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("user_id", orderInfo.UserId),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// Generate a Subscribe Token through orderNo
|
||||
// find user subscribe by subscribe token
|
||||
userSub, err := l.svc.UserModel.FindOneSubscribeByToken(ctx, orderInfo.SubscribeToken)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find user subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("order_id", orderInfo.Id),
|
||||
)
|
||||
return err
|
||||
}
|
||||
userSub.Download = 0
|
||||
userSub.Upload = 0
|
||||
userSub.Status = 1
|
||||
// update user subscribe
|
||||
err = l.svc.UserModel.UpdateSubscribe(ctx, userSub)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Update user subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
sub, err := l.svc.SubscribeModel.FindOne(ctx, userSub.SubscribeId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find subscribe failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("subscribe_id", userSub.SubscribeId),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
userTelegramChatId, ok := findTelegram(userInfo)
|
||||
// SendMessage To Telegram
|
||||
if ok {
|
||||
text, err := tool.RenderTemplateToString(telegram.ResetTrafficNotify, map[string]string{
|
||||
"OrderNo": orderInfo.OrderNo,
|
||||
"SubscribeName": sub.Name,
|
||||
"ResetTime": time.Now().Format("2006-01-02 15:04:05"),
|
||||
"ExpireTime": userSub.ExpireTime.Format("2006-01-02 15:04:05"),
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendUserNotifyWithTelegram(userTelegramChatId, text)
|
||||
}
|
||||
|
||||
// send message to admin
|
||||
text, err := tool.RenderTemplateToString(telegram.AdminOrderNotify, map[string]string{
|
||||
"OrderNo": orderInfo.OrderNo,
|
||||
"TradeNo": orderInfo.TradeNo,
|
||||
"SubscribeName": "流量重置",
|
||||
"OrderAmount": fmt.Sprintf("%.2f", float64(orderInfo.Price)/100),
|
||||
"OrderStatus": "已支付",
|
||||
"OrderTime": orderInfo.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||
"PaymentMethod": orderInfo.Method,
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render AdminOrderNotify template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendAdminNotifyWithTelegram(ctx, text)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Recharge Recharge to user
|
||||
func (l *ActivateOrderLogic) Recharge(ctx context.Context, orderInfo *order.Order) error {
|
||||
// find user by user id
|
||||
userInfo, err := l.svc.UserModel.FindOne(ctx, orderInfo.UserId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Find user failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("user_id", orderInfo.UserId),
|
||||
)
|
||||
return err
|
||||
}
|
||||
userInfo.Balance += orderInfo.Price
|
||||
// update user
|
||||
err = l.svc.DB.Transaction(func(tx *gorm.DB) error {
|
||||
err = l.svc.UserModel.Update(ctx, userInfo, tx)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Update user failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
// Create Balance Log
|
||||
balanceLog := user.BalanceLog{
|
||||
UserId: orderInfo.UserId,
|
||||
Amount: orderInfo.Price,
|
||||
Type: 1,
|
||||
OrderId: orderInfo.Id,
|
||||
Balance: userInfo.Balance,
|
||||
}
|
||||
err = l.svc.UserModel.InsertBalanceLog(ctx, &balanceLog, tx)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Insert balance log failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Database transaction failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
userTelegramChatId, ok := findTelegram(userInfo)
|
||||
// SendMessage To Telegram
|
||||
if ok {
|
||||
text, err := tool.RenderTemplateToString(telegram.RechargeNotify, map[string]string{
|
||||
"OrderAmount": fmt.Sprintf("%.2f", float64(orderInfo.Price)/100),
|
||||
"PaymentMethod": orderInfo.Method,
|
||||
"Time": orderInfo.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||
"Balance": fmt.Sprintf("%.2f", float64(userInfo.Balance)/100),
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendUserNotifyWithTelegram(userTelegramChatId, text)
|
||||
}
|
||||
// send message to admin
|
||||
text, err := tool.RenderTemplateToString(telegram.AdminOrderNotify, map[string]string{
|
||||
"OrderNo": orderInfo.OrderNo,
|
||||
"TradeNo": orderInfo.TradeNo,
|
||||
"OrderAmount": fmt.Sprintf("%.2f", float64(orderInfo.Price)/100),
|
||||
"SubscribeName": "余额充值",
|
||||
"OrderStatus": "已支付",
|
||||
"OrderTime": orderInfo.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||
"PaymentMethod": orderInfo.Method,
|
||||
})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Render AdminOrderNotify template failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
l.sendAdminNotifyWithTelegram(ctx, text)
|
||||
return nil
|
||||
}
|
||||
|
||||
// sendUserNotifyWithTelegram send message to user
|
||||
func (l *ActivateOrderLogic) sendUserNotifyWithTelegram(chatId int64, text string) {
|
||||
msg := tgbotapi.NewMessage(chatId, text)
|
||||
msg.ParseMode = "markdown"
|
||||
_, err := l.svc.TelegramBot.Send(msg)
|
||||
if err != nil {
|
||||
logger.Error("[ActivateOrderLogic] Send telegram user message failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// sendAdminNotifyWithTelegram send message to admin
|
||||
func (l *ActivateOrderLogic) sendAdminNotifyWithTelegram(ctx context.Context, text string) {
|
||||
admins, err := l.svc.UserModel.QueryAdminUsers(ctx)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Query admin users failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return
|
||||
}
|
||||
for _, admin := range admins {
|
||||
telegramId, ok := findTelegram(admin)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
msg := tgbotapi.NewMessage(telegramId, text)
|
||||
msg.ParseMode = "markdown"
|
||||
_, err := l.svc.TelegramBot.Send(msg)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] Send telegram admin message failed",
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// findTelegram find user telegram id
|
||||
func findTelegram(u *user.User) (int64, bool) {
|
||||
for _, item := range u.AuthMethods {
|
||||
if item.AuthType == "telegram" {
|
||||
// string to int64
|
||||
parseInt, err := strconv.ParseInt(item.AuthIdentifier, 10, 64)
|
||||
if err != nil {
|
||||
return 0, false
|
||||
}
|
||||
return parseInt, true
|
||||
}
|
||||
|
||||
}
|
||||
return 0, false
|
||||
}
|
||||
@@ -0,0 +1,161 @@
|
||||
package orderLogic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"github.com/hibiken/asynq"
|
||||
order2 "github.com/perfect-panel/ppanel-server/internal/logic/public/order"
|
||||
"github.com/perfect-panel/ppanel-server/internal/model/payment"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/payment/alipay"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/payment/payssion"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/payment/stripe"
|
||||
"github.com/perfect-panel/ppanel-server/queue/types"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
type CheckOrderLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewCheckOrderLogic(svc *svc.ServiceContext) *CheckOrderLogic {
|
||||
return &CheckOrderLogic{
|
||||
svc: svc,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *CheckOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task) error {
|
||||
|
||||
orderList, err := l.svc.OrderModel.QueryPendingOrders(ctx)
|
||||
if err != nil {
|
||||
logger.Errorf("query pending orders error: %v", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
logger.Infof("查到订单数据: %v", orderList)
|
||||
for _, order := range orderList {
|
||||
paymentConfig, err := l.svc.PaymentModel.FindOne(ctx, order.PaymentId)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckOrder] Find payment config failed", logger.Field("error", err.Error()), logger.Field("paymentMark", order.Method))
|
||||
continue
|
||||
}
|
||||
logger.Infof("查到配置数据[%s]: %v", order.Method, orderList)
|
||||
var flag bool
|
||||
switch order.Method {
|
||||
case order2.AlipayF2f:
|
||||
if l.queryAlipay(paymentConfig, order.TradeNo) {
|
||||
flag = true
|
||||
}
|
||||
break
|
||||
case order2.Payssion:
|
||||
logger.Infof("匹配配置类型: %v", order2.Payssion)
|
||||
if l.queryPayssion(paymentConfig, order.OrderNo) {
|
||||
flag = true
|
||||
}
|
||||
break
|
||||
case order2.StripeWeChatPay:
|
||||
if l.queryStripe(paymentConfig, order.TradeNo) {
|
||||
flag = true
|
||||
}
|
||||
break
|
||||
default:
|
||||
logger.Infow("[CheckOrder] Unsupported payment method", logger.Field("paymentMethod", order.Method))
|
||||
continue
|
||||
}
|
||||
logger.Infof("[CheckOrder] Unsupported payment method[%v]", flag)
|
||||
if flag {
|
||||
err := l.svc.OrderModel.UpdateOrderStatus(ctx, order.OrderNo, 2)
|
||||
if err != nil {
|
||||
logger.Errorf("[CheckOrder] query order status error: %v", zap.Error(err))
|
||||
}
|
||||
logger.Info("[CheckOrder] Notify status success", logger.Field("orderNo", order.TradeNo))
|
||||
payload := types.ForthwithActivateOrderPayload{
|
||||
OrderNo: order.OrderNo,
|
||||
}
|
||||
bytes, err := json.Marshal(&payload)
|
||||
if err != nil {
|
||||
logger.Error("[CheckOrder] Marshal payload failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
task := asynq.NewTask(types.ForthwithActivateOrder, bytes)
|
||||
taskInfo, err := l.svc.Queue.EnqueueContext(ctx, task)
|
||||
if err != nil {
|
||||
logger.Error("[CheckOrder] Enqueue task failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
logger.Info("[CheckOrder] Enqueue task success", logger.Field("taskInfo", taskInfo))
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// queryAlipay Query Alipay payment status
|
||||
//
|
||||
//nolint:unused
|
||||
func (l *CheckOrderLogic) queryAlipay(paymentConfig *payment.Payment, TradeNo string) bool {
|
||||
config := payment.AlipayF2FConfig{}
|
||||
if err := json.Unmarshal([]byte(paymentConfig.Config), &config); err != nil {
|
||||
zap.S().Errorw("[CheckOrder] Unmarshal payment config failed", logger.Field("error", err.Error()), logger.Field("config", paymentConfig.Config))
|
||||
return false
|
||||
}
|
||||
client := alipay.NewClient(alipay.Config{
|
||||
AppId: config.AppId,
|
||||
PrivateKey: config.PrivateKey,
|
||||
PublicKey: config.PublicKey,
|
||||
InvoiceName: config.InvoiceName,
|
||||
})
|
||||
status, err := client.QueryTrade(context.Background(), TradeNo)
|
||||
if err != nil {
|
||||
zap.S().Errorw("[CheckOrder] Query trade failed", logger.Field("error", err.Error()), logger.Field("TradeNo", TradeNo))
|
||||
return false
|
||||
}
|
||||
if status == alipay.Success || status == alipay.Finished {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// queryStripe Query Stripe payment status
|
||||
//
|
||||
//nolint:unused
|
||||
func (l *CheckOrderLogic) queryStripe(paymentConfig *payment.Payment, TradeNo string) bool {
|
||||
config := payment.StripeConfig{}
|
||||
if err := json.Unmarshal([]byte(paymentConfig.Config), &config); err != nil {
|
||||
zap.S().Errorw("[CheckOrder] Unmarshal payment config failed", logger.Field("error", err.Error()), logger.Field("config", paymentConfig.Config))
|
||||
return false
|
||||
}
|
||||
client := stripe.NewClient(stripe.Config{
|
||||
PublicKey: config.PublicKey,
|
||||
SecretKey: config.SecretKey,
|
||||
WebhookSecret: config.WebhookSecret,
|
||||
})
|
||||
status, err := client.QueryOrderStatus(TradeNo)
|
||||
if err != nil {
|
||||
zap.S().Errorw("[CheckOrder] Query order status failed", logger.Field("error", err.Error()), logger.Field("TradeNo", TradeNo))
|
||||
return false
|
||||
}
|
||||
return status
|
||||
}
|
||||
|
||||
// queryPayssion Query Stripe payment status
|
||||
//
|
||||
//nolint:unused
|
||||
func (l *CheckOrderLogic) queryPayssion(paymentConfig *payment.Payment, TradeNo string) bool {
|
||||
zap.S().Infof("[CheckOrder]1 Query Payssion called")
|
||||
payssionConfig := payment.PayssionConfig{}
|
||||
if err := json.Unmarshal([]byte(paymentConfig.Config), &payssionConfig); err != nil {
|
||||
zap.S().Errorw("[CheckOrder] Unmarshal error", logger.Field("error", err.Error()))
|
||||
return false
|
||||
}
|
||||
zap.S().Infof("[CheckOrder]2 Query Payssion called")
|
||||
client := payssion.NewClient(payssionConfig.ApiKey, payssionConfig.SecretKey, payssionConfig.PmId, payssionConfig.Currency, payssionConfig.QueryUrl, payssionConfig.CreateUrl)
|
||||
// create payment
|
||||
result, err := client.QueryOrder(TradeNo)
|
||||
if err != nil {
|
||||
zap.S().Errorw("[CheckOrder] Query order status failed", logger.Field("error", err.Error()), logger.Field("TradeNo", TradeNo))
|
||||
return false
|
||||
}
|
||||
zap.S().Infof("[CheckOrder]3 Query Payssion called")
|
||||
return result.Transaction.State == "completed"
|
||||
}
|
||||
@@ -0,0 +1,47 @@
|
||||
package orderLogic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/perfect-panel/ppanel-server/internal/logic/public/order"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
internal "github.com/perfect-panel/ppanel-server/internal/types"
|
||||
"github.com/perfect-panel/ppanel-server/queue/types"
|
||||
)
|
||||
|
||||
type DeferCloseOrderLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewDeferCloseOrderLogic(svc *svc.ServiceContext) *DeferCloseOrderLogic {
|
||||
return &DeferCloseOrderLogic{
|
||||
svc: svc,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *DeferCloseOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task) error {
|
||||
payload := types.DeferCloseOrderPayload{}
|
||||
if err := json.Unmarshal(task.Payload(), &payload); err != nil {
|
||||
logger.WithContext(ctx).Error("[DeferCloseOrderLogic] Unmarshal payload failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("payload", string(task.Payload())),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
|
||||
err := order.NewCloseOrderLogic(ctx, l.svc).CloseOrder(&internal.CloseOrderRequest{
|
||||
OrderNo: payload.OrderNo,
|
||||
})
|
||||
count, ok := asynq.GetRetryCount(ctx)
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
if err != nil && count < 3 {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,73 @@
|
||||
package smslogic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/perfect-panel/ppanel-server/internal/model/log"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/constant"
|
||||
"github.com/perfect-panel/ppanel-server/pkg/sms"
|
||||
"github.com/perfect-panel/ppanel-server/queue/types"
|
||||
)
|
||||
|
||||
type SmsSendCount struct {
|
||||
Count int `json:"count"`
|
||||
CreateAt int64 `json:"create_at"`
|
||||
}
|
||||
|
||||
type SendSmsLogic struct {
|
||||
svcCtx *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewSendSmsLogic(svcCtx *svc.ServiceContext) *SendSmsLogic {
|
||||
return &SendSmsLogic{
|
||||
svcCtx: svcCtx,
|
||||
}
|
||||
}
|
||||
func (l *SendSmsLogic) ProcessTask(ctx context.Context, task *asynq.Task) error {
|
||||
var payload types.SendSmsPayload
|
||||
if err := json.Unmarshal(task.Payload(), &payload); err != nil {
|
||||
logger.WithContext(ctx).Error("[SendSmsLogic] Unmarshal payload failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("payload", task.Payload()),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
client, err := sms.NewSender(l.svcCtx.Config.Mobile.Platform, l.svcCtx.Config.Mobile.PlatformConfig)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[SendSmsLogic] New send sms client failed", logger.Field("error", err.Error()), logger.Field("payload", payload))
|
||||
return err
|
||||
}
|
||||
createSms := &log.MessageLog{
|
||||
Type: log.Mobile.String(),
|
||||
Platform: l.svcCtx.Config.Mobile.Platform,
|
||||
To: fmt.Sprintf("+%s%s", payload.TelephoneArea, payload.Telephone),
|
||||
Subject: constant.ParseVerifyType(payload.Type).String(),
|
||||
Content: "",
|
||||
}
|
||||
err = client.SendCode(payload.TelephoneArea, payload.Telephone, payload.Content)
|
||||
|
||||
createSms.Content = client.GetSendCodeContent(payload.Content)
|
||||
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[SendSmsLogic] Send sms failed", logger.Field("error", err.Error()), logger.Field("payload", payload))
|
||||
if l.svcCtx.Config.Model != constant.DevMode {
|
||||
createSms.Status = 2
|
||||
} else {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
createSms.Status = 1
|
||||
logger.WithContext(ctx).Info("[SendSmsLogic] Send sms", logger.Field("telephone", payload.Telephone), logger.Field("content", createSms.Content))
|
||||
err = l.svcCtx.LogModel.InsertMessageLog(ctx, createSms)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[SendSmsLogic] Send sms failed", logger.Field("error", err.Error()), logger.Field("payload", payload))
|
||||
return nil
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,215 @@
|
||||
package subscription
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"text/template"
|
||||
"time"
|
||||
|
||||
queue "github.com/perfect-panel/ppanel-server/queue/types"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/perfect-panel/ppanel-server/internal/model/user"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type CheckSubscriptionLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewCheckSubscriptionLogic(svc *svc.ServiceContext) *CheckSubscriptionLogic {
|
||||
return &CheckSubscriptionLogic{
|
||||
svc: svc,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *CheckSubscriptionLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error {
|
||||
logger.Infof("[CheckSubscription] Start check subscription: %s", time.Now().Format("2006-01-02 15:04:05"))
|
||||
// Check subscription traffic
|
||||
err := l.svc.UserModel.Transaction(ctx, func(db *gorm.DB) error {
|
||||
var list []*user.Subscribe
|
||||
err := db.Model(&user.Subscribe{}).Where("upload + download >= traffic AND status = 1 AND traffic > 0 ").Find(&list).Error
|
||||
if err != nil {
|
||||
logger.Errorw("[Check Subscription Traffic] Query subscribe failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
var ids []int64
|
||||
for _, item := range list {
|
||||
ids = append(ids, item.Id)
|
||||
}
|
||||
if len(ids) > 0 {
|
||||
err = db.Model(&user.Subscribe{}).Where("id IN ?", ids).Updates(map[string]interface{}{
|
||||
"status": 2,
|
||||
"finished_at": time.Now(),
|
||||
}).Error
|
||||
if err != nil {
|
||||
logger.Errorw("[Check Subscription Traffic] Update subscribe status failed", logger.Field("error", err.Error()))
|
||||
return nil
|
||||
}
|
||||
err = l.sendTrafficNotify(ctx, ids)
|
||||
if err != nil {
|
||||
logger.Errorw("[Check Subscription Traffic] Send email failed", logger.Field("error", err.Error()))
|
||||
return nil
|
||||
}
|
||||
|
||||
if len(list) > 0 {
|
||||
if err = l.svc.UserModel.ClearSubscribeCache(ctx, list...); err != nil {
|
||||
logger.Errorw("[Check Subscription Traffic] Clear subscribe cache failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
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")
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
logger.Error("[CheckSubscription] Transaction failed", logger.Field("error", err.Error()))
|
||||
}
|
||||
// Check subscription expire
|
||||
err = l.svc.UserModel.Transaction(ctx, func(db *gorm.DB) error {
|
||||
var list []*user.Subscribe
|
||||
err = db.Model(&user.Subscribe{}).Where("`status` = 1 AND `expire_time` < ? AND `expire_time` != ? and `finished_at` IS NULL", time.Now(), time.UnixMilli(0)).Find(&list).Error
|
||||
if err != nil {
|
||||
logger.Error("[Check Subscription] Find subscribe failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
var ids []int64
|
||||
for _, item := range list {
|
||||
ids = append(ids, item.Id)
|
||||
}
|
||||
if len(ids) > 0 {
|
||||
err = db.Model(&user.Subscribe{}).Where("id IN ?", ids).Update("status", 3).Error
|
||||
if err != nil {
|
||||
logger.Error("[Check Subscription Expire] Update subscribe status failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
err = l.sendExpiredNotify(ctx, ids)
|
||||
if err != nil {
|
||||
logger.Error("[Check Subscription Expire] Send email failed", logger.Field("error", err.Error()))
|
||||
return nil
|
||||
}
|
||||
if len(list) > 0 {
|
||||
if err = l.svc.UserModel.ClearSubscribeCache(ctx, list...); err != nil {
|
||||
logger.Errorw("[Check Subscription Traffic] Clear subscribe cache failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
}
|
||||
logger.Info("[Check Subscription Expire] Update subscribe status", logger.Field("user_ids", ids), logger.Field("count", int64(len(ids))))
|
||||
} else {
|
||||
logger.Info("[Check Subscription Expire] No subscribe need to update")
|
||||
}
|
||||
return l.svc.UserModel.ClearSubscribeCache(ctx, list...)
|
||||
})
|
||||
if err != nil {
|
||||
logger.Info("[CheckSubscription] Transaction failed", logger.Field("error", err.Error()))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (l *CheckSubscriptionLogic) sendExpiredNotify(ctx context.Context, subs []int64) error {
|
||||
for _, id := range subs {
|
||||
sub, err := l.svc.UserModel.FindOneUserSubscribe(ctx, id)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] FindOneUserSubscribe failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
method, err := l.svc.UserModel.FindUserAuthMethodByUserId(ctx, "email", sub.UserId)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] FindUserAuthMethodByUserId failed", logger.Field("error", err.Error()), logger.Field("user_id", sub.UserId))
|
||||
continue
|
||||
}
|
||||
var taskPayload queue.SendEmailPayload
|
||||
taskPayload.Email = method.AuthIdentifier
|
||||
taskPayload.Subject = "Subscription Expired"
|
||||
tpl, err := template.New("Expired").Parse(l.svc.Config.Email.ExpirationEmailTemplate)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Parse template failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
var result bytes.Buffer
|
||||
err = tpl.Execute(&result, map[string]interface{}{
|
||||
"SiteLogo": l.svc.Config.Site.SiteLogo,
|
||||
"SiteName": l.svc.Config.Site.SiteName,
|
||||
"ExpireDate": sub.ExpireTime.Format("2006-01-02 15:04:05"),
|
||||
})
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Execute template failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
taskPayload.Content = result.String()
|
||||
payloadBuy, err := json.Marshal(taskPayload)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Marshal payload failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
task := asynq.NewTask(queue.ForthwithSendEmail, payloadBuy, asynq.MaxRetry(3))
|
||||
taskInfo, err := l.svc.Queue.Enqueue(task)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Enqueue task failed", logger.Field("error", err.Error()), logger.Field("payload", string(payloadBuy)))
|
||||
continue
|
||||
}
|
||||
logger.Infow("[CheckSubscription] Send email success",
|
||||
logger.Field("taskID", taskInfo.ID), logger.Field("User", sub.UserId),
|
||||
logger.Field("Email", method.AuthIdentifier),
|
||||
)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (l *CheckSubscriptionLogic) sendTrafficNotify(ctx context.Context, subs []int64) error {
|
||||
for _, id := range subs {
|
||||
sub, err := l.svc.UserModel.FindOneUserSubscribe(ctx, id)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] FindOneUserSubscribe failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
method, err := l.svc.UserModel.FindUserAuthMethodByUserId(ctx, "email", sub.UserId)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] FindUserAuthMethodByUserId failed", logger.Field("error", err.Error()), logger.Field("user_id", sub.UserId))
|
||||
continue
|
||||
}
|
||||
var taskPayload queue.SendEmailPayload
|
||||
taskPayload.Email = method.AuthIdentifier
|
||||
taskPayload.Subject = "Subscription Traffic Exceed"
|
||||
tpl, err := template.New("Traffic").Parse(l.svc.Config.Email.TrafficExceedEmailTemplate)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Parse template failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
var result bytes.Buffer
|
||||
err = tpl.Execute(&result, map[string]interface{}{
|
||||
"SiteLogo": l.svc.Config.Site.SiteLogo,
|
||||
"SiteName": l.svc.Config.Site.SiteName,
|
||||
})
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Execute template failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
taskPayload.Content = result.String()
|
||||
payloadBuy, err := json.Marshal(taskPayload)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Marshal payload failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
task := asynq.NewTask(queue.ForthwithSendEmail, payloadBuy, asynq.MaxRetry(3))
|
||||
taskInfo, err := l.svc.Queue.Enqueue(task)
|
||||
if err != nil {
|
||||
logger.Errorw("[CheckSubscription] Enqueue task failed", logger.Field("error", err.Error()), logger.Field("payload", string(payloadBuy)))
|
||||
continue
|
||||
}
|
||||
logger.Infow("[CheckSubscription] Send email success",
|
||||
logger.Field("taskID", taskInfo.ID), logger.Field("User", sub.UserId),
|
||||
logger.Field("Email", method.AuthIdentifier),
|
||||
)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,166 @@
|
||||
package traffic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/perfect-panel/ppanel-server/internal/config"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"github.com/perfect-panel/ppanel-server/internal/types"
|
||||
)
|
||||
|
||||
type ServerDataLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewServerDataLogic(svc *svc.ServiceContext) *ServerDataLogic {
|
||||
return &ServerDataLogic{
|
||||
svc: svc,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *ServerDataLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error {
|
||||
serverData := types.ServerTotalDataResponse{}
|
||||
|
||||
top10ServerToday, top10ServerYesterday, top10UserToday, top10UserYesterday := l.getRanking(ctx)
|
||||
if len(top10ServerToday) == 0 {
|
||||
top10ServerToday = make([]types.ServerTrafficData, 0)
|
||||
}
|
||||
if len(top10ServerYesterday) == 0 {
|
||||
top10ServerYesterday = make([]types.ServerTrafficData, 0)
|
||||
}
|
||||
if len(top10UserToday) == 0 {
|
||||
top10UserToday = make([]types.UserTrafficData, 0)
|
||||
}
|
||||
if len(top10UserYesterday) == 0 {
|
||||
top10UserYesterday = make([]types.UserTrafficData, 0)
|
||||
}
|
||||
serverData.ServerTrafficRankingToday = top10ServerToday
|
||||
serverData.ServerTrafficRankingYesterday = top10ServerYesterday
|
||||
serverData.UserTrafficRankingToday = top10UserToday
|
||||
serverData.UserTrafficRankingYesterday = top10UserYesterday
|
||||
totalUploadToday, totalDownloadToday, totalDownloadMonthly, totalUploadMonthly := l.trafficCount(ctx)
|
||||
serverData.TodayUpload = totalUploadToday
|
||||
serverData.TodayDownload = totalDownloadToday
|
||||
serverData.MonthlyUpload = totalUploadMonthly
|
||||
serverData.MonthlyDownload = totalDownloadMonthly
|
||||
serverData.UpdatedAt = time.Now().UnixMilli()
|
||||
data, err := json.Marshal(serverData)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Marshal server data failed", logger.Field("error", err.Error()), logger.Field("data", serverData))
|
||||
return err
|
||||
}
|
||||
if err := l.svc.Redis.Set(ctx, config.ServerCountCacheKey, data, -1).Err(); err != nil {
|
||||
logger.Error("[ServerDataLogic] Set server data failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
logger.Info("[ServerDataLogic] Update server data success")
|
||||
return nil
|
||||
}
|
||||
|
||||
func (l *ServerDataLogic) getRanking(ctx context.Context) (top10ServerToday, top10ServerYesterday []types.ServerTrafficData, top10UserToday, top10UserYesterday []types.UserTrafficData) {
|
||||
now := time.Now()
|
||||
// 获取服务器流量排行榜
|
||||
serverToday, err := l.svc.TrafficLogModel.TopServersTrafficByDay(ctx, now, 10)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Get top servers traffic by day failed", logger.Field("error", err.Error()))
|
||||
} else {
|
||||
for _, s := range serverToday {
|
||||
if s.ServerId == 0 {
|
||||
continue
|
||||
}
|
||||
serverInfo, err := l.svc.ServerModel.FindOne(ctx, s.ServerId)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Find server failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
top10ServerToday = append(top10ServerToday, types.ServerTrafficData{
|
||||
ServerId: s.ServerId,
|
||||
Name: serverInfo.Name,
|
||||
Upload: s.Upload,
|
||||
Download: s.Download,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
serverYesterday, err := l.svc.TrafficLogModel.TopServersTrafficByDay(ctx, now.AddDate(0, 0, -1), 10)
|
||||
if err != nil {
|
||||
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)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Find server failed", logger.Field("error", err.Error()))
|
||||
continue
|
||||
}
|
||||
top10ServerYesterday = append(top10ServerYesterday, types.ServerTrafficData{
|
||||
ServerId: s.ServerId,
|
||||
Name: serverInfo.Name,
|
||||
Upload: s.Upload,
|
||||
Download: s.Download,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// 获取用户流量排行榜
|
||||
userToday, err := l.svc.TrafficLogModel.TopUsersTrafficByDay(ctx, now, 10)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Get top users traffic by day failed", logger.Field("error", err.Error()))
|
||||
} else {
|
||||
for _, u := range userToday {
|
||||
//userInfo, err := l.svc.UserModel.FindOne(ctx, u.UserId)
|
||||
//if err != nil {
|
||||
// logx.Error("[ServerDataLogic] Find user failed", logx.Field("error", err.Error()))
|
||||
// continue
|
||||
//}
|
||||
top10UserToday = append(top10UserToday, types.UserTrafficData{
|
||||
SID: u.UserId,
|
||||
Upload: u.Upload,
|
||||
Download: u.Download,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
userYesterday, err := l.svc.TrafficLogModel.TopUsersTrafficByDay(ctx, now.AddDate(0, 0, -1), 10)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Get top users traffic by day failed", logger.Field("error", err.Error()))
|
||||
} else {
|
||||
for _, u := range userYesterday {
|
||||
//userInfo, err := l.svc.UserModel.FindOne(ctx, u.UserId)
|
||||
//if err != nil {
|
||||
// logx.Error("[ServerDataLogic] Find user failed", logx.Field("error", err.Error()))
|
||||
// continue
|
||||
//}
|
||||
top10UserYesterday = append(top10UserYesterday, types.UserTrafficData{
|
||||
SID: u.UserId,
|
||||
Upload: u.Upload,
|
||||
Download: u.Download,
|
||||
})
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (l *ServerDataLogic) trafficCount(ctx context.Context) (totalUploadToday, totalDownloadToday, totalDownloadMonthly, totalUploadMonthly int64) {
|
||||
now := time.Now()
|
||||
today, err := l.svc.TrafficLogModel.QueryTrafficByDay(ctx, now)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Query traffic by day failed", logger.Field("error", err.Error()))
|
||||
} else {
|
||||
totalUploadToday = today.Upload
|
||||
totalDownloadToday = today.Download
|
||||
}
|
||||
|
||||
monthly, err := l.svc.TrafficLogModel.QueryTrafficByMonthly(ctx, now)
|
||||
if err != nil {
|
||||
logger.Error("[ServerDataLogic] Query traffic by monthly failed", logger.Field("error", err.Error()))
|
||||
} else {
|
||||
totalUploadMonthly = monthly.Upload
|
||||
totalDownloadMonthly = monthly.Download
|
||||
}
|
||||
return
|
||||
}
|
||||
@@ -0,0 +1,98 @@
|
||||
package traffic
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/ppanel-server/pkg/logger"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/perfect-panel/ppanel-server/internal/model/traffic"
|
||||
"github.com/perfect-panel/ppanel-server/internal/svc"
|
||||
"github.com/perfect-panel/ppanel-server/queue/types"
|
||||
)
|
||||
|
||||
//goland:noinspection GoNameStartsWithPackageName
|
||||
type TrafficStatisticsLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewTrafficStatisticsLogic(svc *svc.ServiceContext) *TrafficStatisticsLogic {
|
||||
return &TrafficStatisticsLogic{
|
||||
svc: svc,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *TrafficStatisticsLogic) ProcessTask(ctx context.Context, task *asynq.Task) error {
|
||||
var payload types.TrafficStatistics
|
||||
if err := json.Unmarshal(task.Payload(), &payload); err != nil {
|
||||
logger.WithContext(ctx).Error("[TrafficStatistics] Unmarshal payload failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("payload", string(task.Payload())),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
if len(payload.Logs) == 0 {
|
||||
logger.WithContext(ctx).Error("[TrafficStatistics] Payload is empty")
|
||||
return nil
|
||||
}
|
||||
// query server info
|
||||
serverInfo, err := l.svc.ServerModel.FindOne(ctx, payload.ServerId)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[TrafficStatistics] Find server info failed",
|
||||
logger.Field("serverId", payload.ServerId),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
if serverInfo.TrafficRatio == 0 {
|
||||
logger.WithContext(ctx).Error("[TrafficStatistics] Server log ratio is 0",
|
||||
logger.Field("serverId", payload.ServerId),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
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)
|
||||
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("error", err.Error()),
|
||||
)
|
||||
continue
|
||||
}
|
||||
// query user Subscribe Info
|
||||
sub, err := l.svc.UserModel.FindOneSubscribe(ctx, log.SID)
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[TrafficStatistics] Find user Subscribe Info failed",
|
||||
logger.Field("uid", log.SID),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
continue
|
||||
}
|
||||
|
||||
// create log log
|
||||
if err := l.svc.TrafficLogModel.Insert(ctx, &traffic.TrafficLog{
|
||||
ServerId: payload.ServerId,
|
||||
SubscribeId: log.SID,
|
||||
UserId: sub.UserId,
|
||||
Upload: u,
|
||||
Download: d,
|
||||
Timestamp: now,
|
||||
}); 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("error", err.Error()),
|
||||
)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user