Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 05ee5b9d2e |
+2
-2
@@ -7,8 +7,8 @@ MYSQL_ROOT_PASSWORD=CHANGE_ME_TO_STRONG_PASSWORD
|
||||
# Grafana 管理员密码
|
||||
GRAFANA_PASSWORD=CHANGE_ME_TO_STRONG_PASSWORD
|
||||
|
||||
# PPanel Server 镜像标签(由 CI/CD 传入不可变 tag,如 git SHA)
|
||||
PPANEL_SERVER_TAG=CHANGE_ME_TO_GIT_SHA
|
||||
# PPanel Server 镜像标签(留空使用 latest)
|
||||
PPANEL_SERVER_TAG=latest
|
||||
|
||||
# AWS 区域(香港)
|
||||
AWS_REGION=ap-east-1
|
||||
|
||||
@@ -149,35 +149,6 @@ type (
|
||||
Total int64 `json:"total"`
|
||||
List []CommissionLog `json:"list"`
|
||||
}
|
||||
OrderRefundLog {
|
||||
OrderId int64 `json:"order_id"`
|
||||
OrderNo string `json:"order_no"`
|
||||
OperatorUserId int64 `json:"operator_user_id"`
|
||||
OperatorAuthIdentifier string `json:"operator_auth_identifier,omitempty"`
|
||||
TargetUserId int64 `json:"target_user_id"`
|
||||
UserSubscribeId int64 `json:"user_subscribe_id"`
|
||||
RefererUserId int64 `json:"referer_user_id,omitempty"`
|
||||
CommissionAmount int64 `json:"commission_amount"`
|
||||
Reason string `json:"reason,omitempty"`
|
||||
OrderStatusBefore uint8 `json:"order_status_before"`
|
||||
OrderStatusAfter uint8 `json:"order_status_after"`
|
||||
SubscribeStatusBefore uint8 `json:"subscribe_status_before"`
|
||||
SubscribeStatusAfter uint8 `json:"subscribe_status_after"`
|
||||
SubscribeExpireBefore int64 `json:"subscribe_expire_before"`
|
||||
SubscribeExpireAfter int64 `json:"subscribe_expire_after"`
|
||||
CommissionBefore int64 `json:"commission_before"`
|
||||
CommissionAfter int64 `json:"commission_after"`
|
||||
Timestamp int64 `json:"timestamp"`
|
||||
}
|
||||
FilterOrderRefundLogRequest {
|
||||
FilterLogParams
|
||||
OrderId int64 `form:"order_id,optional"`
|
||||
UserId int64 `form:"user_id,optional"`
|
||||
}
|
||||
FilterOrderRefundLogResponse {
|
||||
Total int64 `json:"total"`
|
||||
List []OrderRefundLog `json:"list"`
|
||||
}
|
||||
GiftLog {
|
||||
Type uint16 `json:"type"`
|
||||
userId int64 `json:"user_id"`
|
||||
@@ -268,30 +239,6 @@ type (
|
||||
OccurredAt int64 `json:"occurred_at"`
|
||||
CreatedAt int64 `json:"created_at"`
|
||||
}
|
||||
GetLogMessageRawRequest {
|
||||
Id int64 `form:"id" validate:"required"`
|
||||
}
|
||||
GetLogMessageRawResponse {
|
||||
Id int64 `json:"id"`
|
||||
Platform string `json:"platform"`
|
||||
AppVersion string `json:"app_version"`
|
||||
OsName string `json:"os_name"`
|
||||
OsVersion string `json:"os_version"`
|
||||
DeviceId string `json:"device_id"`
|
||||
UserId *int64 `json:"user_id"`
|
||||
SessionId string `json:"session_id"`
|
||||
Level uint8 `json:"level"`
|
||||
ErrorCode string `json:"error_code"`
|
||||
Message string `json:"message"`
|
||||
Stack string `json:"stack"`
|
||||
Context interface{} `json:"context"`
|
||||
ClientIP string `json:"client_ip"`
|
||||
UserAgent string `json:"user_agent"`
|
||||
Locale string `json:"locale"`
|
||||
Digest string `json:"digest"`
|
||||
OccurredAt int64 `json:"occurred_at"`
|
||||
CreatedAt int64 `json:"created_at"`
|
||||
}
|
||||
)
|
||||
|
||||
@server (
|
||||
@@ -344,10 +291,6 @@ service ppanel {
|
||||
@handler FilterCommissionLog
|
||||
get /commission/list (FilterCommissionLogRequest) returns (FilterCommissionLogResponse)
|
||||
|
||||
@doc "Filter order refund log"
|
||||
@handler FilterOrderRefundLog
|
||||
get /order/refund/list (FilterOrderRefundLogRequest) returns (FilterOrderRefundLogResponse)
|
||||
|
||||
@doc "Filter gift log"
|
||||
@handler FilterGiftLog
|
||||
get /gift/list (FilterGiftLogRequest) returns (FilterGiftLogResponse)
|
||||
@@ -371,9 +314,5 @@ service ppanel {
|
||||
@doc "Get error log message detail"
|
||||
@handler GetErrorLogMessageDetail
|
||||
get /error_message/detail returns (GetErrorLogMessageDetailResponse)
|
||||
|
||||
@doc "Get log message raw detail (temporary)"
|
||||
@handler GetLogMessageRaw
|
||||
get /message/detail (GetLogMessageRawRequest) returns (GetLogMessageRawResponse)
|
||||
}
|
||||
|
||||
|
||||
@@ -33,10 +33,6 @@ type (
|
||||
PaymentId int64 `json:"payment_id,omitempty"`
|
||||
TradeNo string `json:"trade_no,omitempty"`
|
||||
}
|
||||
RefundOrderRequest {
|
||||
Id int64 `json:"id" validate:"required"`
|
||||
Reason string `json:"reason,omitempty" validate:"omitempty,max=500"`
|
||||
}
|
||||
ActivateOrderRequest {
|
||||
OrderNo string `json:"order_no" validate:"required"`
|
||||
}
|
||||
@@ -72,10 +68,6 @@ service ppanel {
|
||||
@handler UpdateOrderStatus
|
||||
put /status (UpdateOrderStatusRequest)
|
||||
|
||||
@doc "Refund order"
|
||||
@handler RefundOrder
|
||||
post /refund (RefundOrderRequest)
|
||||
|
||||
@doc "Manually activate order"
|
||||
@handler ActivateOrder
|
||||
post /activate (ActivateOrderRequest)
|
||||
|
||||
@@ -427,7 +427,6 @@ type (
|
||||
FeeAmount int64 `json:"fee_amount"`
|
||||
TradeNo string `json:"trade_no"`
|
||||
Status uint8 `json:"status"`
|
||||
StatusName string `json:"status_name,omitempty"`
|
||||
SubscribeId int64 `json:"subscribe_id"`
|
||||
CreatedAt int64 `json:"created_at"`
|
||||
UpdatedAt int64 `json:"updated_at"`
|
||||
@@ -450,7 +449,6 @@ type (
|
||||
FeeAmount int64 `json:"fee_amount"`
|
||||
TradeNo string `json:"trade_no"`
|
||||
Status uint8 `json:"status"`
|
||||
StatusName string `json:"status_name,omitempty"`
|
||||
SubscribeId int64 `json:"subscribe_id"`
|
||||
Subscribe Subscribe `json:"subscribe"`
|
||||
CreatedAt int64 `json:"created_at"`
|
||||
|
||||
@@ -19,10 +19,9 @@ services:
|
||||
# ----------------------------------------------------
|
||||
# 1. 业务后端 (PPanel Server)
|
||||
# host 网络:可出外网,直接访问 AWS RDS/Redis;通过 127.0.0.1 访问 Tempo
|
||||
# PPANEL_SERVER_TAG 由 CI/CD 传入不可变镜像标签(如 git SHA)
|
||||
# ----------------------------------------------------
|
||||
ppanel-server:
|
||||
image: registry.kxsw.us/vpn-server:${PPANEL_SERVER_TAG:?please set PPANEL_SERVER_TAG to an immutable image tag}
|
||||
image: registry.kxsw.us/vpn-server:${PPANEL_SERVER_TAG:-latest}
|
||||
container_name: ppanel-server
|
||||
restart: always
|
||||
volumes:
|
||||
@@ -120,7 +119,7 @@ services:
|
||||
# 或配置 Nginx 反代(建议加认证)
|
||||
# ----------------------------------------------------
|
||||
grafana:
|
||||
image: grafana/grafana:13.0.1
|
||||
image: grafana/grafana:latest
|
||||
container_name: ppanel-grafana
|
||||
restart: always
|
||||
ports:
|
||||
@@ -155,7 +154,7 @@ services:
|
||||
# 6. Prometheus (指标采集)
|
||||
# ----------------------------------------------------
|
||||
prometheus:
|
||||
image: prom/prometheus:v3.11.3
|
||||
image: prom/prometheus:latest
|
||||
container_name: ppanel-prometheus
|
||||
restart: always
|
||||
ports:
|
||||
@@ -180,7 +179,7 @@ services:
|
||||
# 7. Nginx Exporter (监控宿主机 Nginx)
|
||||
# ----------------------------------------------------
|
||||
nginx-exporter:
|
||||
image: nginx/nginx-prometheus-exporter:1.5.0
|
||||
image: nginx/nginx-prometheus-exporter:latest
|
||||
container_name: ppanel-nginx-exporter
|
||||
restart: always
|
||||
command:
|
||||
@@ -199,7 +198,7 @@ services:
|
||||
# 8. Node Exporter (宿主机监控)
|
||||
# ----------------------------------------------------
|
||||
node-exporter:
|
||||
image: prom/node-exporter:v1.11.1
|
||||
image: prom/node-exporter:latest
|
||||
container_name: ppanel-node-exporter
|
||||
restart: always
|
||||
volumes:
|
||||
@@ -222,7 +221,7 @@ services:
|
||||
# 9. cAdvisor (容器监控)
|
||||
# ----------------------------------------------------
|
||||
cadvisor:
|
||||
image: gcr.io/cadvisor/cadvisor:v0.55.1
|
||||
image: gcr.io/cadvisor/cadvisor:latest
|
||||
container_name: ppanel-cadvisor
|
||||
restart: always
|
||||
volumes:
|
||||
|
||||
@@ -1,13 +0,0 @@
|
||||
-- Remove app_account_token column from order table if it exists
|
||||
SET @col_exists = (SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'order' AND COLUMN_NAME = 'app_account_token');
|
||||
SET @sql = IF(@col_exists = 1, 'ALTER TABLE `order` DROP COLUMN `app_account_token`', 'SELECT 1');
|
||||
PREPARE stmt FROM @sql;
|
||||
EXECUTE stmt;
|
||||
DEALLOCATE PREPARE stmt;
|
||||
|
||||
-- Remove subscription_user_id column from order table if it exists
|
||||
SET @col_exists2 = (SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'order' AND COLUMN_NAME = 'subscription_user_id');
|
||||
SET @sql2 = IF(@col_exists2 = 1, 'ALTER TABLE `order` DROP COLUMN `subscription_user_id`', 'SELECT 1');
|
||||
PREPARE stmt2 FROM @sql2;
|
||||
EXECUTE stmt2;
|
||||
DEALLOCATE PREPARE stmt2;
|
||||
@@ -1,19 +0,0 @@
|
||||
-- Rollback: re-deduct commission for users with pending (status=0) withdrawals.
|
||||
-- This re-applies the OLD behaviour where commission is deducted on application.
|
||||
-- Only run this if you are rolling back to the old code; do NOT run against
|
||||
-- the new code or commission will be double-deducted on approval.
|
||||
|
||||
UPDATE `user` u
|
||||
JOIN (
|
||||
SELECT user_id, COALESCE(SUM(amount), 0) AS pending_total
|
||||
FROM withdrawals
|
||||
WHERE status = 0
|
||||
GROUP BY user_id
|
||||
) p ON u.id = p.user_id
|
||||
SET u.commission = u.commission - p.pending_total
|
||||
WHERE p.pending_total > 0;
|
||||
|
||||
-- Remove the migration log entries written by the up migration.
|
||||
DELETE FROM system_logs
|
||||
WHERE type = 33
|
||||
AND content LIKE '%migration: refund pending withdrawal commission (HIF-22)%';
|
||||
@@ -1,65 +0,0 @@
|
||||
-- Migration: refund commission for existing pending (status=0) withdrawals
|
||||
--
|
||||
-- Under the old logic, commission was deducted when a withdrawal was submitted.
|
||||
-- Under the new logic, commission is only deducted on approval.
|
||||
-- This migration refunds the deducted amounts back to each user so that
|
||||
-- the system is in a consistent state before the new code is deployed.
|
||||
--
|
||||
-- Idempotency: the UPDATE only touches rows whose commission would need
|
||||
-- to increase, and each execution produces the same result because
|
||||
-- COALESCE(SUM(amount),0) is deterministic given the same pending set.
|
||||
-- Running this script multiple times is safe only if no new pending
|
||||
-- withdrawals are created between runs; deploy new code immediately after.
|
||||
|
||||
-- Compatibility: some historical databases missed migration 02122, so the
|
||||
-- withdrawals table may not exist yet. Create it idempotently before the
|
||||
-- refund logic so this migration can self-heal older installations.
|
||||
CREATE TABLE IF NOT EXISTS `withdrawals` (
|
||||
`id` BIGINT NOT NULL AUTO_INCREMENT COMMENT 'Primary Key',
|
||||
`user_id` BIGINT NOT NULL COMMENT 'User ID',
|
||||
`amount` BIGINT NOT NULL COMMENT 'Withdrawal Amount',
|
||||
`content` TEXT COMMENT 'Withdrawal Content',
|
||||
`status` TINYINT(1) NOT NULL DEFAULT 0 COMMENT 'Withdrawal Status',
|
||||
`reason` VARCHAR(500) NOT NULL DEFAULT '' COMMENT 'Rejection Reason',
|
||||
`created_at` DATETIME NOT NULL COMMENT 'Creation Time',
|
||||
`updated_at` DATETIME NOT NULL COMMENT 'Update Time',
|
||||
PRIMARY KEY (`id`),
|
||||
KEY `idx_user_id` (`user_id`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
|
||||
|
||||
INSERT IGNORE INTO `system` (`category`, `key`, `value`, `type`, `desc`, `created_at`, `updated_at`)
|
||||
VALUES
|
||||
('invite', 'WithdrawalMethod', '', 'string', 'withdrawal method', '2025-04-22 14:25:16.637', '2025-04-22 14:25:16.637');
|
||||
|
||||
-- Step 1: refund commission for all users with pending withdrawals.
|
||||
UPDATE `user` u
|
||||
JOIN (
|
||||
SELECT user_id, COALESCE(SUM(amount), 0) AS pending_total
|
||||
FROM withdrawals
|
||||
WHERE status = 0
|
||||
GROUP BY user_id
|
||||
) p ON u.id = p.user_id
|
||||
SET u.commission = u.commission + p.pending_total
|
||||
WHERE p.pending_total > 0;
|
||||
|
||||
-- Step 2: write a migration log entry for each refunded user.
|
||||
INSERT INTO system_logs (type, date, object_id, content, created_at)
|
||||
SELECT
|
||||
33 AS type,
|
||||
DATE(NOW()) AS date,
|
||||
p.user_id AS object_id,
|
||||
JSON_OBJECT(
|
||||
'type', 99,
|
||||
'amount', p.pending_total,
|
||||
'order_no', '',
|
||||
'timestamp', UNIX_TIMESTAMP(NOW()) * 1000,
|
||||
'note', 'migration: refund pending withdrawal commission (HIF-22)'
|
||||
) AS content,
|
||||
NOW() AS created_at
|
||||
FROM (
|
||||
SELECT user_id, COALESCE(SUM(amount), 0) AS pending_total
|
||||
FROM withdrawals
|
||||
WHERE status = 0
|
||||
GROUP BY user_id
|
||||
HAVING pending_total > 0
|
||||
) p;
|
||||
@@ -1,2 +0,0 @@
|
||||
-- Remove activation_context column from order table
|
||||
ALTER TABLE `order` DROP COLUMN IF EXISTS `activation_context`;
|
||||
@@ -1,6 +0,0 @@
|
||||
-- Add activation_context column to order table for Redis fallback persistence (idempotent)
|
||||
SET @col_exists = (SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = 'order' AND COLUMN_NAME = 'activation_context');
|
||||
SET @sql = IF(@col_exists = 0, 'ALTER TABLE `order` ADD COLUMN `activation_context` TEXT DEFAULT NULL COMMENT ''Activation context JSON (guest/redemption info for DB fallback)'' AFTER `app_account_token`', 'SELECT 1');
|
||||
PREPARE stmt FROM @sql;
|
||||
EXECUTE stmt;
|
||||
DEALLOCATE PREPARE stmt;
|
||||
@@ -1,25 +0,0 @@
|
||||
package log
|
||||
|
||||
import (
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/perfect-panel/server/internal/logic/admin/log"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
"github.com/perfect-panel/server/pkg/result"
|
||||
)
|
||||
|
||||
func FilterOrderRefundLogHandler(svcCtx *svc.ServiceContext) func(c *gin.Context) {
|
||||
return func(c *gin.Context) {
|
||||
var req types.FilterOrderRefundLogRequest
|
||||
_ = c.ShouldBind(&req)
|
||||
validateErr := svcCtx.Validate(&req)
|
||||
if validateErr != nil {
|
||||
result.ParamErrorResult(c, validateErr)
|
||||
return
|
||||
}
|
||||
|
||||
l := log.NewFilterOrderRefundLogLogic(c.Request.Context(), svcCtx)
|
||||
resp, err := l.FilterOrderRefundLog(&req)
|
||||
result.HttpResult(c, resp, err)
|
||||
}
|
||||
}
|
||||
@@ -1,23 +0,0 @@
|
||||
package log
|
||||
|
||||
import (
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/perfect-panel/server/internal/logic/admin/log"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
"github.com/perfect-panel/server/pkg/result"
|
||||
)
|
||||
|
||||
func GetLogMessageRawHandler(svcCtx *svc.ServiceContext) func(c *gin.Context) {
|
||||
return func(c *gin.Context) {
|
||||
var req types.GetLogMessageRawRequest
|
||||
_ = c.ShouldBind(&req)
|
||||
if err := svcCtx.Validate(&req); err != nil {
|
||||
result.ParamErrorResult(c, err)
|
||||
return
|
||||
}
|
||||
l := log.NewGetLogMessageRawLogic(c.Request.Context(), svcCtx)
|
||||
resp, err := l.GetLogMessageRaw(&req)
|
||||
result.HttpResult(c, resp, err)
|
||||
}
|
||||
}
|
||||
@@ -1,25 +0,0 @@
|
||||
package order
|
||||
|
||||
import (
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/perfect-panel/server/internal/logic/admin/order"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
"github.com/perfect-panel/server/pkg/result"
|
||||
)
|
||||
|
||||
func RefundOrderHandler(svcCtx *svc.ServiceContext) func(c *gin.Context) {
|
||||
return func(c *gin.Context) {
|
||||
var req types.RefundOrderRequest
|
||||
_ = c.ShouldBind(&req)
|
||||
validateErr := svcCtx.Validate(&req)
|
||||
if validateErr != nil {
|
||||
result.ParamErrorResult(c, validateErr)
|
||||
return
|
||||
}
|
||||
|
||||
l := order.NewRefundOrderLogic(c.Request.Context(), svcCtx)
|
||||
err := l.RefundOrder(&req)
|
||||
result.HttpResult(c, nil, err)
|
||||
}
|
||||
}
|
||||
@@ -250,18 +250,12 @@ func RegisterHandlers(router *gin.Engine, serverCtx *svc.ServiceContext) {
|
||||
// Filter commission log
|
||||
adminLogGroupRouter.GET("/commission/list", adminLog.FilterCommissionLogHandler(serverCtx))
|
||||
|
||||
// Filter order refund log
|
||||
adminLogGroupRouter.GET("/order/refund/list", adminLog.FilterOrderRefundLogHandler(serverCtx))
|
||||
|
||||
// Filter email log
|
||||
adminLogGroupRouter.GET("/email/list", adminLog.FilterEmailLogHandler(serverCtx))
|
||||
|
||||
// Get error log message detail
|
||||
adminLogGroupRouter.GET("/error_message/detail", adminLog.GetErrorLogMessageDetailHandler(serverCtx))
|
||||
|
||||
// Get log message raw detail (temporary)
|
||||
adminLogGroupRouter.GET("/message/detail", adminLog.GetLogMessageRawHandler(serverCtx))
|
||||
|
||||
// Get error log message list
|
||||
adminLogGroupRouter.GET("/error_message/list", adminLog.GetErrorLogMessageListHandler(serverCtx))
|
||||
|
||||
@@ -347,9 +341,6 @@ func RegisterHandlers(router *gin.Engine, serverCtx *svc.ServiceContext) {
|
||||
// Update order status
|
||||
adminOrderGroupRouter.PUT("/status", adminOrder.UpdateOrderStatusHandler(serverCtx))
|
||||
|
||||
// Refund order
|
||||
adminOrderGroupRouter.POST("/refund", adminOrder.RefundOrderHandler(serverCtx))
|
||||
|
||||
// Manually activate order
|
||||
adminOrderGroupRouter.POST("/activate", adminOrder.ActivateOrderHandler(serverCtx))
|
||||
}
|
||||
|
||||
@@ -29,7 +29,7 @@ func (l *DeleteSubscribeApplicationLogic) DeleteSubscribeApplication(req *types.
|
||||
err := l.svcCtx.ClientModel.Delete(l.ctx, req.Id)
|
||||
if err != nil {
|
||||
l.Errorf("Failed to delete subscribe application with ID %d: %v", req.Id, err)
|
||||
return errors.Wrap(xerr.NewErrCode(xerr.DatabaseDeletedError), err.Error())
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseDeletedError), err.Error())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,77 +0,0 @@
|
||||
package log
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
modellog "github.com/perfect-panel/server/internal/model/log"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
"github.com/perfect-panel/server/pkg/logger"
|
||||
"github.com/perfect-panel/server/pkg/xerr"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
type FilterOrderRefundLogLogic struct {
|
||||
logger.Logger
|
||||
ctx context.Context
|
||||
svcCtx *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewFilterOrderRefundLogLogic(ctx context.Context, svcCtx *svc.ServiceContext) *FilterOrderRefundLogLogic {
|
||||
return &FilterOrderRefundLogLogic{
|
||||
Logger: logger.WithContext(ctx),
|
||||
ctx: ctx,
|
||||
svcCtx: svcCtx,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *FilterOrderRefundLogLogic) FilterOrderRefundLog(req *types.FilterOrderRefundLogRequest) (resp *types.FilterOrderRefundLogResponse, err error) {
|
||||
data, total, err := l.svcCtx.LogModel.FilterSystemLog(l.ctx, &modellog.FilterParams{
|
||||
Page: req.Page,
|
||||
Size: req.Size,
|
||||
Data: req.Date,
|
||||
Search: req.Search,
|
||||
Type: modellog.TypeOrderRefund.Uint8(),
|
||||
ObjectID: req.OrderId,
|
||||
})
|
||||
if err != nil {
|
||||
l.Errorw("Query order refund log failed", logger.Field("error", err.Error()))
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "Query order refund log failed")
|
||||
}
|
||||
|
||||
var list []types.OrderRefundLog
|
||||
for _, item := range data {
|
||||
var content modellog.OrderRefund
|
||||
if err := content.Unmarshal([]byte(item.Content)); err != nil {
|
||||
l.Errorf("unmarshal order refund log content failed: %v", err.Error())
|
||||
continue
|
||||
}
|
||||
if req.UserId > 0 && content.TargetUserId != req.UserId && content.RefererUserId != req.UserId && content.OperatorUserId != req.UserId {
|
||||
continue
|
||||
}
|
||||
list = append(list, types.OrderRefundLog{
|
||||
OrderId: content.OrderId,
|
||||
OrderNo: content.OrderNo,
|
||||
OperatorUserId: content.OperatorUserId,
|
||||
OperatorAuthIdentifier: content.OperatorAuthIdentifier,
|
||||
TargetUserId: content.TargetUserId,
|
||||
UserSubscribeId: content.UserSubscribeId,
|
||||
RefererUserId: content.RefererUserId,
|
||||
CommissionAmount: content.CommissionAmount,
|
||||
Reason: content.Reason,
|
||||
OrderStatusBefore: content.OrderStatusBefore,
|
||||
OrderStatusAfter: content.OrderStatusAfter,
|
||||
SubscribeStatusBefore: content.SubscribeStatusBefore,
|
||||
SubscribeStatusAfter: content.SubscribeStatusAfter,
|
||||
SubscribeExpireBefore: content.SubscribeExpireBefore,
|
||||
SubscribeExpireAfter: content.SubscribeExpireAfter,
|
||||
CommissionBefore: content.CommissionBefore,
|
||||
CommissionAfter: content.CommissionAfter,
|
||||
Timestamp: content.Timestamp,
|
||||
})
|
||||
}
|
||||
return &types.FilterOrderRefundLogResponse{
|
||||
Total: total,
|
||||
List: list,
|
||||
}, nil
|
||||
}
|
||||
@@ -1,66 +0,0 @@
|
||||
package log
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
"github.com/perfect-panel/server/pkg/logger"
|
||||
"github.com/perfect-panel/server/pkg/xerr"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
type GetLogMessageRawLogic struct {
|
||||
logger.Logger
|
||||
ctx context.Context
|
||||
svcCtx *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewGetLogMessageRawLogic(ctx context.Context, svcCtx *svc.ServiceContext) *GetLogMessageRawLogic {
|
||||
return &GetLogMessageRawLogic{Logger: logger.WithContext(ctx), ctx: ctx, svcCtx: svcCtx}
|
||||
}
|
||||
|
||||
func (l *GetLogMessageRawLogic) GetLogMessageRaw(req *types.GetLogMessageRawRequest) (resp *types.GetLogMessageRawResponse, err error) {
|
||||
row, err := l.svcCtx.LogMessageModel.FindOne(l.ctx, req.Id)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "FindOne log_message id=%d: %v", req.Id, err)
|
||||
}
|
||||
|
||||
var contextJSON json.RawMessage
|
||||
if row.Context != "" {
|
||||
if json.Valid([]byte(row.Context)) {
|
||||
contextJSON = json.RawMessage(row.Context)
|
||||
} else {
|
||||
b, _ := json.Marshal(row.Context)
|
||||
contextJSON = b
|
||||
}
|
||||
}
|
||||
|
||||
var occurredAt int64
|
||||
if row.OccurredAt != nil {
|
||||
occurredAt = row.OccurredAt.UnixMilli()
|
||||
}
|
||||
|
||||
return &types.GetLogMessageRawResponse{
|
||||
Id: row.Id,
|
||||
Platform: row.Platform,
|
||||
AppVersion: row.AppVersion,
|
||||
OsName: row.OsName,
|
||||
OsVersion: row.OsVersion,
|
||||
DeviceId: row.DeviceId,
|
||||
UserId: row.UserId,
|
||||
SessionId: row.SessionId,
|
||||
Level: row.Level,
|
||||
ErrorCode: row.ErrorCode,
|
||||
Message: row.Message,
|
||||
Stack: row.Stack,
|
||||
Context: contextJSON,
|
||||
ClientIP: row.ClientIP,
|
||||
UserAgent: row.UserAgent,
|
||||
Locale: row.Locale,
|
||||
Digest: row.Digest,
|
||||
OccurredAt: occurredAt,
|
||||
CreatedAt: row.CreatedAt.UnixMilli(),
|
||||
}, nil
|
||||
}
|
||||
@@ -35,28 +35,6 @@ func (l *GetOrderListLogic) GetOrderList(req *types.GetOrderListRequest) (resp *
|
||||
resp = &types.GetOrderListResponse{}
|
||||
resp.List = make([]types.Order, 0)
|
||||
tool.DeepCopy(&resp.List, list)
|
||||
for i := range resp.List {
|
||||
resp.List[i].StatusName = orderStatusName(resp.List[i].Status)
|
||||
}
|
||||
resp.Total = total
|
||||
return
|
||||
}
|
||||
|
||||
func orderStatusName(status uint8) string {
|
||||
switch status {
|
||||
case 1:
|
||||
return "pending"
|
||||
case 2:
|
||||
return "paid"
|
||||
case 3:
|
||||
return "closed"
|
||||
case 4:
|
||||
return "failed"
|
||||
case 5:
|
||||
return "finished"
|
||||
case 6:
|
||||
return "refunded"
|
||||
default:
|
||||
return "unknown"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,298 +0,0 @@
|
||||
package order
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/server/internal/model/log"
|
||||
modelorder "github.com/perfect-panel/server/internal/model/order"
|
||||
modeluser "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/constant"
|
||||
"github.com/perfect-panel/server/pkg/logger"
|
||||
"github.com/perfect-panel/server/pkg/xerr"
|
||||
"github.com/pkg/errors"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
const (
|
||||
orderStatusRefunded = 6
|
||||
)
|
||||
|
||||
type RefundOrderLogic struct {
|
||||
logger.Logger
|
||||
ctx context.Context
|
||||
svcCtx *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewRefundOrderLogic(ctx context.Context, svcCtx *svc.ServiceContext) *RefundOrderLogic {
|
||||
return &RefundOrderLogic{
|
||||
Logger: logger.WithContext(ctx),
|
||||
ctx: ctx,
|
||||
svcCtx: svcCtx,
|
||||
}
|
||||
}
|
||||
|
||||
func (l *RefundOrderLogic) RefundOrder(req *types.RefundOrderRequest) error {
|
||||
operator, ok := l.ctx.Value(constant.CtxKeyUser).(*modeluser.User)
|
||||
if !ok {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "Invalid Access")
|
||||
}
|
||||
|
||||
reason := strings.TrimSpace(req.Reason)
|
||||
var cachesToClear []*modeluser.Subscribe
|
||||
var planCacheIDs []int64
|
||||
var userCacheTargets []*modeluser.User
|
||||
|
||||
err := l.svcCtx.DB.WithContext(l.ctx).Transaction(func(tx *gorm.DB) error {
|
||||
var orderInfo modelorder.Order
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Model(&modelorder.Order{}).
|
||||
Where("id = ?", req.Id).
|
||||
First(&orderInfo).Error; err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.OrderNotExist), "order %d not found", req.Id)
|
||||
}
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query order failed: %v", err)
|
||||
}
|
||||
|
||||
if orderInfo.Status == orderStatusRefunded {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.OrderAlreadyRefunded), "order %d already refunded", orderInfo.Id)
|
||||
}
|
||||
if orderInfo.Status != 2 && orderInfo.Status != 5 {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.OrderStatusError), "order %d status %d is not refundable", orderInfo.Id, orderInfo.Status)
|
||||
}
|
||||
|
||||
userSub, err := l.lockRefundTargetSubscription(tx, &orderInfo)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
orderStatusBefore := orderInfo.Status
|
||||
subStatusBefore := userSub.Status
|
||||
subExpireBefore := userSub.ExpireTime
|
||||
now := time.Now()
|
||||
|
||||
referer, commissionAmount, err := l.lockCommissionSource(tx, orderInfo.OrderNo, orderInfo.Commission)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
var commissionBefore int64
|
||||
var commissionAfter int64
|
||||
if referer != nil {
|
||||
commissionBefore = referer.Commission
|
||||
commissionAfter = referer.Commission - commissionAmount
|
||||
if err := tx.Model(&modeluser.User{}).
|
||||
Where("id = ?", referer.Id).
|
||||
UpdateColumn("commission", gorm.Expr("commission - ?", commissionAmount)).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "refund commission failed: %v", err)
|
||||
}
|
||||
|
||||
commissionLog := log.Commission{
|
||||
Type: log.CommissionTypeRefund,
|
||||
Amount: -commissionAmount,
|
||||
OrderNo: orderInfo.OrderNo,
|
||||
Timestamp: now.UnixMilli(),
|
||||
}
|
||||
content, marshalErr := commissionLog.Marshal()
|
||||
if marshalErr != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.ERROR), "marshal commission refund log failed: %v", marshalErr)
|
||||
}
|
||||
if err := tx.Model(&log.SystemLog{}).Create(&log.SystemLog{
|
||||
Type: log.TypeCommission.Uint8(),
|
||||
Date: now.Format(time.DateOnly),
|
||||
ObjectID: referer.Id,
|
||||
Content: string(content),
|
||||
CreatedAt: now,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "insert commission refund log failed: %v", err)
|
||||
}
|
||||
referer.Commission = commissionAfter
|
||||
userCacheTargets = append(userCacheTargets, referer)
|
||||
}
|
||||
|
||||
if err := tx.Model(&modelorder.Order{}).
|
||||
Where("id = ? AND status IN ?", orderInfo.Id, []int{2, 5}).
|
||||
Update("status", orderStatusRefunded).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "refund order status update failed: %v", err)
|
||||
}
|
||||
orderInfo.Status = orderStatusRefunded
|
||||
|
||||
userSub.Status = 3
|
||||
userSub.ExpireTime = now.Add(-time.Second)
|
||||
userSub.FinishedAt = &now
|
||||
if err := tx.Model(&modeluser.Subscribe{}).
|
||||
Where("id = ?", userSub.Id).
|
||||
Updates(map[string]interface{}{
|
||||
"status": userSub.Status,
|
||||
"expire_time": userSub.ExpireTime,
|
||||
"finished_at": userSub.FinishedAt,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "update subscription failed: %v", err)
|
||||
}
|
||||
|
||||
refundLog, err := l.buildRefundAuditLog(operator, &orderInfo, userSub, referer, commissionAmount, reason, orderStatusBefore, subStatusBefore, subExpireBefore, commissionBefore, commissionAfter, now)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
logContent, err := refundLog.Marshal()
|
||||
if err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.ERROR), "marshal refund audit log failed: %v", err)
|
||||
}
|
||||
if err := tx.Model(&log.SystemLog{}).Create(&log.SystemLog{
|
||||
Type: log.TypeOrderRefund.Uint8(),
|
||||
Date: now.Format(time.DateOnly),
|
||||
ObjectID: orderInfo.Id,
|
||||
Content: string(logContent),
|
||||
CreatedAt: now,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "insert refund audit log failed: %v", err)
|
||||
}
|
||||
|
||||
cachesToClear = append(cachesToClear, userSub)
|
||||
if userSub.SubscribeId > 0 {
|
||||
planCacheIDs = append(planCacheIDs, userSub.SubscribeId)
|
||||
}
|
||||
if orderInfo.UserId > 0 {
|
||||
orderUser := &modeluser.User{Id: orderInfo.UserId}
|
||||
userCacheTargets = append(userCacheTargets, orderUser)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
l.Errorw("[RefundOrder] refund failed", logger.Field("error", err.Error()), logger.Field("order_id", req.Id))
|
||||
return err
|
||||
}
|
||||
|
||||
if len(cachesToClear) > 0 {
|
||||
if clearErr := l.svcCtx.UserModel.ClearSubscribeCache(l.ctx, cachesToClear...); clearErr != nil {
|
||||
l.Errorw("[RefundOrder] clear subscribe cache failed", logger.Field("error", clearErr.Error()), logger.Field("order_id", req.Id))
|
||||
}
|
||||
}
|
||||
for _, subscribeID := range planCacheIDs {
|
||||
if clearErr := l.svcCtx.SubscribeModel.ClearCache(l.ctx, subscribeID); clearErr != nil {
|
||||
l.Errorw("[RefundOrder] clear subscribe cache failed", logger.Field("error", clearErr.Error()), logger.Field("subscribe_id", subscribeID))
|
||||
}
|
||||
}
|
||||
if len(userCacheTargets) > 0 {
|
||||
if clearErr := l.svcCtx.UserModel.ClearUserCache(l.ctx, userCacheTargets...); clearErr != nil {
|
||||
l.Errorw("[RefundOrder] clear user cache failed", logger.Field("error", clearErr.Error()))
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (l *RefundOrderLogic) lockRefundTargetSubscription(tx *gorm.DB, orderInfo *modelorder.Order) (*modeluser.Subscribe, error) {
|
||||
var userSub modeluser.Subscribe
|
||||
err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Model(&modeluser.Subscribe{}).
|
||||
Where("order_id = ?", orderInfo.Id).
|
||||
First(&userSub).Error
|
||||
if err == nil {
|
||||
return &userSub, nil
|
||||
}
|
||||
if !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query refund subscription failed: %v", err)
|
||||
}
|
||||
|
||||
if orderInfo.Type != 2 {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.OrderRefundNoSubscription), "order %d has no linked subscription", orderInfo.Id)
|
||||
}
|
||||
|
||||
if orderInfo.ParentId == 0 {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.OrderRefundNoSubscription), "renewal order %d has no parent order", orderInfo.Id)
|
||||
}
|
||||
|
||||
err = tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Model(&modeluser.Subscribe{}).
|
||||
Where("order_id = ?", orderInfo.ParentId).
|
||||
First(&userSub).Error
|
||||
if err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.OrderRefundNoSubscription), "renewal order %d parent subscription not found", orderInfo.Id)
|
||||
}
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query renewal subscription failed: %v", err)
|
||||
}
|
||||
return &userSub, nil
|
||||
}
|
||||
|
||||
func (l *RefundOrderLogic) lockCommissionSource(tx *gorm.DB, orderNo string, orderCommission int64) (*modeluser.User, int64, error) {
|
||||
var commissionLogs []log.SystemLog
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Model(&log.SystemLog{}).
|
||||
Where("type = ? AND content LIKE ?", log.TypeCommission.Uint8(), fmt.Sprintf("%%\"order_no\":\"%s\"%%", orderNo)).
|
||||
Order("id DESC").
|
||||
Find(&commissionLogs).Error; err != nil {
|
||||
return nil, 0, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query commission log failed: %v", err)
|
||||
}
|
||||
|
||||
for _, item := range commissionLogs {
|
||||
var content log.Commission
|
||||
if err := content.Unmarshal([]byte(item.Content)); err != nil {
|
||||
continue
|
||||
}
|
||||
if content.Type != log.CommissionTypePurchase && content.Type != log.CommissionTypeRenewal {
|
||||
continue
|
||||
}
|
||||
|
||||
var referer modeluser.User
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Model(&modeluser.User{}).
|
||||
Where("id = ?", item.ObjectID).
|
||||
First(&referer).Error; err != nil {
|
||||
return nil, 0, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query commission owner failed: %v", err)
|
||||
}
|
||||
return &referer, content.Amount, nil
|
||||
}
|
||||
|
||||
if orderCommission > 0 {
|
||||
return nil, 0, errors.Wrapf(xerr.NewErrCode(xerr.OrderRefundCommissionMismatch), "order %s commission owner not found", orderNo)
|
||||
}
|
||||
return nil, 0, nil
|
||||
}
|
||||
|
||||
func (l *RefundOrderLogic) buildRefundAuditLog(
|
||||
operator *modeluser.User,
|
||||
orderInfo *modelorder.Order,
|
||||
userSub *modeluser.Subscribe,
|
||||
referer *modeluser.User,
|
||||
commissionAmount int64,
|
||||
reason string,
|
||||
orderStatusBefore uint8,
|
||||
subStatusBefore uint8,
|
||||
subExpireBefore time.Time,
|
||||
commissionBefore int64,
|
||||
commissionAfter int64,
|
||||
now time.Time,
|
||||
) (*log.OrderRefund, error) {
|
||||
refundLog := &log.OrderRefund{
|
||||
OrderId: orderInfo.Id,
|
||||
OrderNo: orderInfo.OrderNo,
|
||||
OperatorUserId: operator.Id,
|
||||
TargetUserId: userSub.UserId,
|
||||
UserSubscribeId: userSub.Id,
|
||||
CommissionAmount: commissionAmount,
|
||||
Reason: reason,
|
||||
OrderStatusBefore: orderStatusBefore,
|
||||
OrderStatusAfter: orderInfo.Status,
|
||||
SubscribeStatusBefore: subStatusBefore,
|
||||
SubscribeStatusAfter: userSub.Status,
|
||||
SubscribeExpireBefore: subExpireBefore.UnixMilli(),
|
||||
SubscribeExpireAfter: userSub.ExpireTime.UnixMilli(),
|
||||
CommissionBefore: commissionBefore,
|
||||
CommissionAfter: commissionAfter,
|
||||
Timestamp: now.UnixMilli(),
|
||||
}
|
||||
if referer != nil {
|
||||
refundLog.RefererUserId = referer.Id
|
||||
}
|
||||
if len(operator.AuthMethods) > 0 {
|
||||
refundLog.OperatorAuthIdentifier = operator.AuthMethods[0].AuthIdentifier
|
||||
}
|
||||
return refundLog, nil
|
||||
}
|
||||
@@ -1,85 +0,0 @@
|
||||
package order
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
modelorder "github.com/perfect-panel/server/internal/model/order"
|
||||
modeluser "github.com/perfect-panel/server/internal/model/user"
|
||||
)
|
||||
|
||||
func TestOrderStatusName(t *testing.T) {
|
||||
tests := map[uint8]string{
|
||||
1: "pending",
|
||||
2: "paid",
|
||||
3: "closed",
|
||||
4: "failed",
|
||||
5: "finished",
|
||||
6: "refunded",
|
||||
7: "unknown",
|
||||
}
|
||||
|
||||
for input, want := range tests {
|
||||
if got := orderStatusName(input); got != want {
|
||||
t.Fatalf("status %d: got %q want %q", input, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildRefundAuditLog(t *testing.T) {
|
||||
logic := &RefundOrderLogic{}
|
||||
now := time.Unix(1710000000, 0)
|
||||
expireBefore := now.Add(24 * time.Hour)
|
||||
expireAfter := now.Add(-time.Second)
|
||||
operator := &modeluser.User{
|
||||
Id: 100,
|
||||
AuthMethods: []modeluser.AuthMethods{
|
||||
{AuthIdentifier: "admin@example.com"},
|
||||
},
|
||||
}
|
||||
orderInfo := &modelorder.Order{Id: 10, OrderNo: "ORD-1", Status: orderStatusRefunded}
|
||||
userSub := &modeluser.Subscribe{Id: 20, UserId: 30, Status: 3, ExpireTime: expireAfter}
|
||||
referer := &modeluser.User{Id: 40}
|
||||
|
||||
got, err := logic.buildRefundAuditLog(
|
||||
operator,
|
||||
orderInfo,
|
||||
userSub,
|
||||
referer,
|
||||
500,
|
||||
"manual refund",
|
||||
5,
|
||||
1,
|
||||
expireBefore,
|
||||
900,
|
||||
400,
|
||||
now,
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatalf("buildRefundAuditLog error: %v", err)
|
||||
}
|
||||
if got.OrderId != 10 || got.OrderNo != "ORD-1" {
|
||||
t.Fatalf("unexpected order info: %+v", got)
|
||||
}
|
||||
if got.OperatorUserId != 100 || got.OperatorAuthIdentifier != "admin@example.com" {
|
||||
t.Fatalf("unexpected operator info: %+v", got)
|
||||
}
|
||||
if got.TargetUserId != 30 || got.UserSubscribeId != 20 {
|
||||
t.Fatalf("unexpected subscription info: %+v", got)
|
||||
}
|
||||
if got.RefererUserId != 40 || got.CommissionAmount != 500 {
|
||||
t.Fatalf("unexpected commission info: %+v", got)
|
||||
}
|
||||
if got.OrderStatusBefore != 5 || got.OrderStatusAfter != orderStatusRefunded {
|
||||
t.Fatalf("unexpected order status: %+v", got)
|
||||
}
|
||||
if got.SubscribeStatusBefore != 1 || got.SubscribeStatusAfter != 3 {
|
||||
t.Fatalf("unexpected subscribe status: %+v", got)
|
||||
}
|
||||
if got.SubscribeExpireBefore != expireBefore.UnixMilli() || got.SubscribeExpireAfter != expireAfter.UnixMilli() {
|
||||
t.Fatalf("unexpected expire transition: %+v", got)
|
||||
}
|
||||
if got.CommissionBefore != 900 || got.CommissionAfter != 400 {
|
||||
t.Fatalf("unexpected commission transition: %+v", got)
|
||||
}
|
||||
}
|
||||
@@ -80,7 +80,7 @@ func (l *ResetSortWithNodeLogic) ResetSortWithNode(req *types.ResetSortRequest)
|
||||
})
|
||||
if err != nil {
|
||||
l.Errorw("[NodeSort] Update Database Error: ", logger.Field("error", err.Error()))
|
||||
return errors.Wrap(xerr.NewErrCode(xerr.DatabaseUpdateError), err.Error())
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), err.Error())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -80,7 +80,7 @@ func (l *ResetSortWithServerLogic) ResetSortWithServer(req *types.ResetSortReque
|
||||
})
|
||||
if err != nil {
|
||||
l.Errorw("[NodeSort] Update Database Error: ", logger.Field("error", err.Error()))
|
||||
return errors.Wrap(xerr.NewErrCode(xerr.DatabaseUpdateError), err.Error())
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), err.Error())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -120,9 +120,7 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
|
||||
if req.ReferCode != "" {
|
||||
userInfo.ReferCode = req.ReferCode
|
||||
}
|
||||
if req.RefererId != 0 {
|
||||
userInfo.RefererId = req.RefererId
|
||||
}
|
||||
userInfo.RefererId = req.RefererId
|
||||
if req.Enable != nil {
|
||||
userInfo.Enable = req.Enable
|
||||
}
|
||||
|
||||
@@ -12,12 +12,10 @@ import (
|
||||
"github.com/perfect-panel/server/pkg/xerr"
|
||||
"github.com/pkg/errors"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
func approveWithdrawal(ctx context.Context, svcCtx *svc.ServiceContext, withdrawalID int64) error {
|
||||
var approvedUserID int64
|
||||
err := svcCtx.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
return svcCtx.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
withdrawal, err := logicCommon.LoadPendingWithdrawalForUpdate(ctx, tx, withdrawalID)
|
||||
if err != nil {
|
||||
if err.Error() == "withdrawal status invalid" {
|
||||
@@ -26,16 +24,6 @@ func approveWithdrawal(ctx context.Context, svcCtx *svc.ServiceContext, withdraw
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "load withdrawal failed: %v", err)
|
||||
}
|
||||
|
||||
// Lock user row and verify sufficient balance before deducting.
|
||||
var u usermodel.User
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Where("id = ?", withdrawal.UserId).First(&u).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "load user failed: %v", err)
|
||||
}
|
||||
if u.Commission < withdrawal.Amount {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.UserCommissionNotEnough), "user %d has insufficient commission balance", withdrawal.UserId)
|
||||
}
|
||||
|
||||
if err := tx.Model(&usermodel.Withdrawal{}).
|
||||
Where("id = ? AND status = 0", withdrawalID).
|
||||
Updates(map[string]interface{}{
|
||||
@@ -45,31 +33,17 @@ func approveWithdrawal(ctx context.Context, svcCtx *svc.ServiceContext, withdraw
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "approve withdrawal failed: %v", err)
|
||||
}
|
||||
|
||||
// Deduct commission atomically inside the transaction.
|
||||
if err := tx.Model(&usermodel.User{}).
|
||||
Where("id = ?", withdrawal.UserId).
|
||||
UpdateColumn("commission", gorm.Expr("commission - ?", withdrawal.Amount)).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "deduct commission failed: %v", err)
|
||||
}
|
||||
|
||||
if err := logicCommon.WriteCommissionLog(tx, withdrawal.UserId, log.CommissionTypeWithdraw, withdrawal.Amount, ""); err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "write commission log failed: %v", err)
|
||||
}
|
||||
|
||||
approvedUserID = withdrawal.UserId
|
||||
return nil
|
||||
})
|
||||
|
||||
if err == nil && approvedUserID > 0 {
|
||||
_ = svcCtx.UserModel.ClearUserCache(ctx, &usermodel.User{Id: approvedUserID})
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func rejectWithdrawal(ctx context.Context, svcCtx *svc.ServiceContext, withdrawalID int64, reason string) error {
|
||||
reason = strings.TrimSpace(reason)
|
||||
return svcCtx.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
_, err := logicCommon.LoadPendingWithdrawalForUpdate(ctx, tx, withdrawalID)
|
||||
withdrawal, err := logicCommon.LoadPendingWithdrawalForUpdate(ctx, tx, withdrawalID)
|
||||
if err != nil {
|
||||
if err.Error() == "withdrawal status invalid" {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.WithdrawalStatusInvalid), "withdrawal %d already processed", withdrawalID)
|
||||
@@ -77,14 +51,23 @@ func rejectWithdrawal(ctx context.Context, svcCtx *svc.ServiceContext, withdrawa
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "load withdrawal failed: %v", err)
|
||||
}
|
||||
|
||||
// Commission was NOT deducted at application time under the new logic,
|
||||
// so rejection requires no refund — only a status update.
|
||||
return tx.Model(&usermodel.Withdrawal{}).
|
||||
if err := tx.Model(&usermodel.Withdrawal{}).
|
||||
Where("id = ? AND status = 0", withdrawalID).
|
||||
Updates(map[string]interface{}{
|
||||
"status": 2,
|
||||
"reason": reason,
|
||||
}).Error
|
||||
}).Error; err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "reject withdrawal failed: %v", err)
|
||||
}
|
||||
|
||||
if err := svcCtx.UserModel.UpdateCommission(ctx, withdrawal.UserId, withdrawal.Amount, tx); err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "refund commission failed: %v", err)
|
||||
}
|
||||
|
||||
if err := logicCommon.WriteCommissionLog(tx, withdrawal.UserId, log.CommissionTypeWithdrawReject, withdrawal.Amount, ""); err != nil {
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "write commission log failed: %v", err)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -1,346 +0,0 @@
|
||||
package common
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
stderrors "errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/server/internal/model/user"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/pkg/xerr"
|
||||
"github.com/pkg/errors"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
const (
|
||||
PromoRuleTypeNewUser = "new_user"
|
||||
PromoRuleTypeInactiveUser = "inactive_user"
|
||||
PromoRuleTypeCampaign = "campaign"
|
||||
|
||||
promoRulesEnabledCacheKey = "promo:rules:enabled"
|
||||
promoSubscribeCachePrefix = "promo:subscribe:"
|
||||
promoCacheTTL = 5 * time.Minute
|
||||
)
|
||||
|
||||
type PromoResult struct {
|
||||
Eligible bool
|
||||
RuleID int64
|
||||
RuleName string
|
||||
RuleType string
|
||||
PromoPrice int64
|
||||
ExpiresAt time.Time
|
||||
}
|
||||
|
||||
type promoRule struct {
|
||||
Id int64 `gorm:"column:id" json:"id"`
|
||||
Name string `gorm:"column:name" json:"name"`
|
||||
Type string `gorm:"column:type" json:"type"`
|
||||
Params string `gorm:"column:params" json:"params"`
|
||||
Priority int64 `gorm:"column:priority" json:"priority"`
|
||||
Enabled bool `gorm:"column:enabled" json:"enabled"`
|
||||
StartTime *time.Time `gorm:"column:start_time" json:"start_time,omitempty"`
|
||||
EndTime *time.Time `gorm:"column:end_time" json:"end_time,omitempty"`
|
||||
DeletedAt gorm.DeletedAt `gorm:"column:deleted_at" json:"deleted_at"`
|
||||
}
|
||||
|
||||
func (promoRule) TableName() string {
|
||||
return "promo_rule"
|
||||
}
|
||||
|
||||
type subscribePromo struct {
|
||||
Id int64 `gorm:"column:id" json:"id"`
|
||||
SubscribeId int64 `gorm:"column:subscribe_id" json:"subscribe_id"`
|
||||
Quantity int64 `gorm:"column:quantity" json:"quantity"`
|
||||
PromoRuleId int64 `gorm:"column:promo_rule_id" json:"promo_rule_id"`
|
||||
PromoPrice int64 `gorm:"column:promo_price" json:"promo_price"`
|
||||
CreatedAt time.Time `gorm:"column:created_at" json:"created_at"`
|
||||
UpdatedAt time.Time `gorm:"column:updated_at" json:"updated_at"`
|
||||
}
|
||||
|
||||
func (subscribePromo) TableName() string {
|
||||
return "subscribe_promo"
|
||||
}
|
||||
|
||||
type promoRuleParams struct {
|
||||
WindowHours int64 `json:"window_hours"`
|
||||
InactiveMonths int `json:"inactive_months"`
|
||||
}
|
||||
|
||||
type promoEligibilitySource interface {
|
||||
UserCreatedAt(ctx context.Context, userID int64) (time.Time, error)
|
||||
LastSubscribeExpireAt(ctx context.Context, userID int64) (time.Time, error)
|
||||
}
|
||||
|
||||
type gormPromoEligibilitySource struct {
|
||||
db *gorm.DB
|
||||
}
|
||||
|
||||
func EvaluatePromo(ctx context.Context, svcCtx *svc.ServiceContext, userID int64, subscribeID int64, quantity int64) (*PromoResult, error) {
|
||||
if svcCtx == nil || svcCtx.DB == nil {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.InvalidParams), "service context is empty")
|
||||
}
|
||||
if userID <= 0 || subscribeID <= 0 || quantity <= 0 {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.InvalidParams), "user id, subscribe id or quantity is empty")
|
||||
}
|
||||
|
||||
rules, err := loadEnabledPromoRules(ctx, svcCtx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(rules) == 0 {
|
||||
return &PromoResult{}, nil
|
||||
}
|
||||
|
||||
prices, err := loadSubscribePromos(ctx, svcCtx, subscribeID, quantity)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(prices) == 0 {
|
||||
return &PromoResult{}, nil
|
||||
}
|
||||
|
||||
return evaluatePromoRulesAt(ctx, rules, prices, &gormPromoEligibilitySource{db: svcCtx.DB}, userID, time.Now())
|
||||
}
|
||||
|
||||
func evaluatePromoRulesAt(
|
||||
ctx context.Context,
|
||||
rules []promoRule,
|
||||
prices []subscribePromo,
|
||||
source promoEligibilitySource,
|
||||
userID int64,
|
||||
now time.Time,
|
||||
) (*PromoResult, error) {
|
||||
priceByRuleID := make(map[int64]int64, len(prices))
|
||||
for _, item := range prices {
|
||||
if item.PromoRuleId <= 0 || item.PromoPrice <= 0 {
|
||||
continue
|
||||
}
|
||||
priceByRuleID[item.PromoRuleId] = item.PromoPrice
|
||||
}
|
||||
|
||||
var userCreatedAt time.Time
|
||||
var userCreatedAtLoaded bool
|
||||
var lastExpireAt time.Time
|
||||
var lastExpireAtLoaded bool
|
||||
|
||||
for _, rule := range rules {
|
||||
promoPrice, ok := priceByRuleID[rule.Id]
|
||||
if !ok || !rule.isAvailableAt(now) {
|
||||
continue
|
||||
}
|
||||
|
||||
params, err := rule.params()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var eligible bool
|
||||
var expiresAt time.Time
|
||||
switch rule.Type {
|
||||
case PromoRuleTypeNewUser:
|
||||
if !userCreatedAtLoaded {
|
||||
userCreatedAt, err = source.UserCreatedAt(ctx, userID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
userCreatedAtLoaded = true
|
||||
}
|
||||
eligible, expiresAt = evaluatePromoNewUser(userCreatedAt, params, now)
|
||||
case PromoRuleTypeInactiveUser:
|
||||
if !lastExpireAtLoaded {
|
||||
lastExpireAt, err = source.LastSubscribeExpireAt(ctx, userID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
lastExpireAtLoaded = true
|
||||
}
|
||||
eligible, expiresAt = evaluatePromoInactiveUser(lastExpireAt, params, rule, now)
|
||||
case PromoRuleTypeCampaign:
|
||||
eligible, expiresAt = evaluatePromoCampaign(rule)
|
||||
default:
|
||||
continue
|
||||
}
|
||||
|
||||
if !eligible {
|
||||
continue
|
||||
}
|
||||
return &PromoResult{
|
||||
Eligible: true,
|
||||
RuleID: rule.Id,
|
||||
RuleName: rule.Name,
|
||||
RuleType: rule.Type,
|
||||
PromoPrice: promoPrice,
|
||||
ExpiresAt: expiresAt,
|
||||
}, nil
|
||||
}
|
||||
|
||||
return &PromoResult{}, nil
|
||||
}
|
||||
|
||||
func evaluatePromoNewUser(userCreatedAt time.Time, params promoRuleParams, now time.Time) (bool, time.Time) {
|
||||
if userCreatedAt.IsZero() || params.WindowHours <= 0 {
|
||||
return false, time.Time{}
|
||||
}
|
||||
expiresAt := userCreatedAt.Add(time.Duration(params.WindowHours) * time.Hour)
|
||||
return now.Before(expiresAt), expiresAt
|
||||
}
|
||||
|
||||
func evaluatePromoInactiveUser(lastExpireAt time.Time, params promoRuleParams, rule promoRule, now time.Time) (bool, time.Time) {
|
||||
if params.InactiveMonths <= 0 {
|
||||
return false, time.Time{}
|
||||
}
|
||||
if lastExpireAt.Equal(time.UnixMilli(0)) || lastExpireAt.After(now) {
|
||||
return false, time.Time{}
|
||||
}
|
||||
if lastExpireAt.IsZero() {
|
||||
return true, rule.expiresAt()
|
||||
}
|
||||
|
||||
threshold := now.AddDate(0, -params.InactiveMonths, 0)
|
||||
return !lastExpireAt.After(threshold), rule.expiresAt()
|
||||
}
|
||||
|
||||
func evaluatePromoCampaign(rule promoRule) (bool, time.Time) {
|
||||
return true, rule.expiresAt()
|
||||
}
|
||||
|
||||
func (r promoRule) isAvailableAt(now time.Time) bool {
|
||||
if r.Id <= 0 || !r.Enabled || r.DeletedAt.Valid {
|
||||
return false
|
||||
}
|
||||
if r.StartTime != nil && now.Before(*r.StartTime) {
|
||||
return false
|
||||
}
|
||||
if r.EndTime != nil && now.After(*r.EndTime) {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func (r promoRule) expiresAt() time.Time {
|
||||
if r.EndTime == nil {
|
||||
return time.Time{}
|
||||
}
|
||||
return *r.EndTime
|
||||
}
|
||||
|
||||
func (r promoRule) params() (promoRuleParams, error) {
|
||||
if r.Params == "" {
|
||||
return promoRuleParams{}, nil
|
||||
}
|
||||
var params promoRuleParams
|
||||
if err := json.Unmarshal([]byte(r.Params), ¶ms); err != nil {
|
||||
return promoRuleParams{}, errors.Wrapf(xerr.NewErrCode(xerr.InvalidParams), "parse promo rule params failed")
|
||||
}
|
||||
return params, nil
|
||||
}
|
||||
|
||||
func (s *gormPromoEligibilitySource) UserCreatedAt(ctx context.Context, userID int64) (time.Time, error) {
|
||||
var item user.User
|
||||
err := s.db.WithContext(ctx).
|
||||
Model(&user.User{}).
|
||||
Where("id = ?", userID).
|
||||
Take(&item).Error
|
||||
if err != nil {
|
||||
return time.Time{}, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query promo user failed")
|
||||
}
|
||||
return item.CreatedAt, nil
|
||||
}
|
||||
|
||||
func (s *gormPromoEligibilitySource) LastSubscribeExpireAt(ctx context.Context, userID int64) (time.Time, error) {
|
||||
var item user.Subscribe
|
||||
err := lastSubscribeExpireQuery(s.db.WithContext(ctx), userID).
|
||||
Take(&item).Error
|
||||
if err != nil {
|
||||
if stderrors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return time.Time{}, nil
|
||||
}
|
||||
return time.Time{}, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query promo user last subscription failed")
|
||||
}
|
||||
return item.ExpireTime, nil
|
||||
}
|
||||
|
||||
func lastSubscribeExpireQuery(db *gorm.DB, userID int64) *gorm.DB {
|
||||
return db.
|
||||
Model(&user.Subscribe{}).
|
||||
Where("user_id = ?", userID).
|
||||
Order(clause.OrderBy{Expression: clause.Expr{
|
||||
SQL: "CASE WHEN expire_time = ? THEN 0 ELSE 1 END ASC, expire_time DESC",
|
||||
Vars: []interface{}{time.UnixMilli(0)},
|
||||
}}).
|
||||
Limit(1)
|
||||
}
|
||||
|
||||
func loadEnabledPromoRules(ctx context.Context, svcCtx *svc.ServiceContext) ([]promoRule, error) {
|
||||
if cached, ok := getPromoCache[[]promoRule](ctx, svcCtx, promoRulesEnabledCacheKey); ok {
|
||||
return cached, nil
|
||||
}
|
||||
|
||||
var rules []promoRule
|
||||
if err := svcCtx.DB.WithContext(ctx).
|
||||
Model(&promoRule{}).
|
||||
Where("enabled = ?", true).
|
||||
Order("priority DESC").
|
||||
Order("id ASC").
|
||||
Find(&rules).Error; err != nil {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query enabled promo rules failed")
|
||||
}
|
||||
|
||||
setPromoCache(ctx, svcCtx, promoRulesEnabledCacheKey, rules)
|
||||
return rules, nil
|
||||
}
|
||||
|
||||
func loadSubscribePromos(ctx context.Context, svcCtx *svc.ServiceContext, subscribeID int64, quantity int64) ([]subscribePromo, error) {
|
||||
cacheKey := fmt.Sprintf("%s%d:%d", promoSubscribeCachePrefix, subscribeID, quantity)
|
||||
if cached, ok := getPromoCache[[]subscribePromo](ctx, svcCtx, cacheKey); ok {
|
||||
return cached, nil
|
||||
}
|
||||
|
||||
var promos []subscribePromo
|
||||
if err := subscribePromoQuery(svcCtx.DB.WithContext(ctx), subscribeID, quantity).
|
||||
Find(&promos).Error; err != nil {
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query subscribe promos failed")
|
||||
}
|
||||
|
||||
setPromoCache(ctx, svcCtx, cacheKey, promos)
|
||||
return promos, nil
|
||||
}
|
||||
|
||||
func subscribePromoQuery(db *gorm.DB, subscribeID int64, quantity int64) *gorm.DB {
|
||||
return db.
|
||||
Model(&subscribePromo{}).
|
||||
Where("subscribe_id = ? AND quantity = ?", subscribeID, quantity)
|
||||
}
|
||||
|
||||
func getPromoCache[T any](ctx context.Context, svcCtx *svc.ServiceContext, key string) (T, bool) {
|
||||
var zero T
|
||||
if svcCtx == nil || svcCtx.Redis == nil {
|
||||
return zero, false
|
||||
}
|
||||
|
||||
value, err := svcCtx.Redis.Get(ctx, key).Result()
|
||||
if err != nil {
|
||||
return zero, false
|
||||
}
|
||||
|
||||
var data T
|
||||
if err = json.Unmarshal([]byte(value), &data); err != nil {
|
||||
_ = svcCtx.Redis.Del(ctx, key).Err()
|
||||
return zero, false
|
||||
}
|
||||
return data, true
|
||||
}
|
||||
|
||||
func setPromoCache(ctx context.Context, svcCtx *svc.ServiceContext, key string, value any) {
|
||||
if svcCtx == nil || svcCtx.Redis == nil {
|
||||
return
|
||||
}
|
||||
payload, err := json.Marshal(value)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_ = svcCtx.Redis.Set(ctx, key, string(payload), promoCacheTTL).Err()
|
||||
}
|
||||
@@ -1,207 +0,0 @@
|
||||
package common
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"gorm.io/driver/mysql"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
)
|
||||
|
||||
type fakePromoEligibilitySource struct {
|
||||
userCreatedAt time.Time
|
||||
lastExpireAt time.Time
|
||||
}
|
||||
|
||||
func (s fakePromoEligibilitySource) UserCreatedAt(context.Context, int64) (time.Time, error) {
|
||||
return s.userCreatedAt, nil
|
||||
}
|
||||
|
||||
func (s fakePromoEligibilitySource) LastSubscribeExpireAt(context.Context, int64) (time.Time, error) {
|
||||
return s.lastExpireAt, nil
|
||||
}
|
||||
|
||||
func TestEvaluatePromoRulesAtPriorityFirstMatch(t *testing.T) {
|
||||
now := time.Date(2026, 5, 27, 12, 0, 0, 0, time.UTC)
|
||||
campaignEnd := now.Add(24 * time.Hour)
|
||||
rules := []promoRule{
|
||||
{
|
||||
Id: 2,
|
||||
Name: "campaign",
|
||||
Type: PromoRuleTypeCampaign,
|
||||
Priority: 20,
|
||||
Enabled: true,
|
||||
EndTime: &campaignEnd,
|
||||
},
|
||||
{
|
||||
Id: 1,
|
||||
Name: "new user",
|
||||
Type: PromoRuleTypeNewUser,
|
||||
Params: `{"window_hours":168}`,
|
||||
Priority: 10,
|
||||
Enabled: true,
|
||||
},
|
||||
}
|
||||
prices := []subscribePromo{
|
||||
{SubscribeId: 100, Quantity: 1, PromoRuleId: 1, PromoPrice: 599},
|
||||
{SubscribeId: 100, Quantity: 1, PromoRuleId: 2, PromoPrice: 499},
|
||||
}
|
||||
|
||||
got, err := evaluatePromoRulesAt(context.Background(), rules, prices, fakePromoEligibilitySource{
|
||||
userCreatedAt: now.Add(-time.Hour),
|
||||
}, 10, now)
|
||||
if err != nil {
|
||||
t.Fatalf("evaluatePromoRulesAt error: %v", err)
|
||||
}
|
||||
if !got.Eligible || got.RuleID != 2 || got.PromoPrice != 499 || got.RuleType != PromoRuleTypeCampaign {
|
||||
t.Fatalf("unexpected promo result: %+v", got)
|
||||
}
|
||||
if !got.ExpiresAt.Equal(campaignEnd) {
|
||||
t.Fatalf("expires_at = %v, want %v", got.ExpiresAt, campaignEnd)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluatePromoRulesAtSkipsUnavailableRules(t *testing.T) {
|
||||
now := time.Date(2026, 5, 27, 12, 0, 0, 0, time.UTC)
|
||||
futureStart := now.Add(time.Hour)
|
||||
expiredEnd := now.Add(-time.Hour)
|
||||
rules := []promoRule{
|
||||
{Id: 1, Type: PromoRuleTypeCampaign, Enabled: true, StartTime: &futureStart},
|
||||
{Id: 2, Type: PromoRuleTypeCampaign, Enabled: true, EndTime: &expiredEnd},
|
||||
{Id: 3, Type: PromoRuleTypeCampaign, Enabled: false},
|
||||
{Id: 4, Type: "unknown", Enabled: true},
|
||||
{Id: 5, Type: PromoRuleTypeCampaign, Enabled: true},
|
||||
}
|
||||
prices := []subscribePromo{
|
||||
{SubscribeId: 100, Quantity: 1, PromoRuleId: 1, PromoPrice: 100},
|
||||
{SubscribeId: 100, Quantity: 1, PromoRuleId: 2, PromoPrice: 100},
|
||||
{SubscribeId: 100, Quantity: 1, PromoRuleId: 3, PromoPrice: 100},
|
||||
{SubscribeId: 100, Quantity: 1, PromoRuleId: 4, PromoPrice: 100},
|
||||
{SubscribeId: 100, Quantity: 1, PromoRuleId: 5, PromoPrice: 88},
|
||||
}
|
||||
|
||||
got, err := evaluatePromoRulesAt(context.Background(), rules, prices, fakePromoEligibilitySource{}, 10, now)
|
||||
if err != nil {
|
||||
t.Fatalf("evaluatePromoRulesAt error: %v", err)
|
||||
}
|
||||
if !got.Eligible || got.RuleID != 5 || got.PromoPrice != 88 {
|
||||
t.Fatalf("unexpected promo result: %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluatePromoNewUser(t *testing.T) {
|
||||
now := time.Date(2026, 5, 27, 12, 0, 0, 0, time.UTC)
|
||||
createdAt := now.Add(-23 * time.Hour)
|
||||
|
||||
eligible, expiresAt := evaluatePromoNewUser(createdAt, promoRuleParams{WindowHours: 24}, now)
|
||||
if !eligible {
|
||||
t.Fatal("new user should be eligible inside configured window")
|
||||
}
|
||||
if !expiresAt.Equal(createdAt.Add(24 * time.Hour)) {
|
||||
t.Fatalf("expires_at = %v, want %v", expiresAt, createdAt.Add(24*time.Hour))
|
||||
}
|
||||
|
||||
eligible, _ = evaluatePromoNewUser(createdAt, promoRuleParams{WindowHours: 12}, now)
|
||||
if eligible {
|
||||
t.Fatal("new user should not be eligible after configured window")
|
||||
}
|
||||
|
||||
eligible, _ = evaluatePromoNewUser(createdAt, promoRuleParams{}, now)
|
||||
if eligible {
|
||||
t.Fatal("new user should not be eligible without positive window_hours")
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluatePromoInactiveUser(t *testing.T) {
|
||||
now := time.Date(2026, 5, 27, 12, 0, 0, 0, time.UTC)
|
||||
endTime := now.Add(48 * time.Hour)
|
||||
rule := promoRule{EndTime: &endTime}
|
||||
|
||||
eligible, expiresAt := evaluatePromoInactiveUser(time.Time{}, promoRuleParams{InactiveMonths: 3}, rule, now)
|
||||
if !eligible {
|
||||
t.Fatal("never purchased user should be eligible for inactive promo")
|
||||
}
|
||||
if !expiresAt.Equal(endTime) {
|
||||
t.Fatalf("expires_at = %v, want %v", expiresAt, endTime)
|
||||
}
|
||||
|
||||
eligible, _ = evaluatePromoInactiveUser(now.AddDate(0, -4, 0), promoRuleParams{InactiveMonths: 3}, rule, now)
|
||||
if !eligible {
|
||||
t.Fatal("expired before inactive threshold should be eligible")
|
||||
}
|
||||
|
||||
eligible, _ = evaluatePromoInactiveUser(now.AddDate(0, -1, 0), promoRuleParams{InactiveMonths: 3}, rule, now)
|
||||
if eligible {
|
||||
t.Fatal("recently expired subscription should not be eligible")
|
||||
}
|
||||
|
||||
eligible, _ = evaluatePromoInactiveUser(time.UnixMilli(0), promoRuleParams{InactiveMonths: 3}, rule, now)
|
||||
if eligible {
|
||||
t.Fatal("unlimited active subscription should not be eligible")
|
||||
}
|
||||
|
||||
eligible, _ = evaluatePromoInactiveUser(now.Add(time.Hour), promoRuleParams{InactiveMonths: 3}, rule, now)
|
||||
if eligible {
|
||||
t.Fatal("currently active subscription should not be eligible")
|
||||
}
|
||||
|
||||
eligible, _ = evaluatePromoInactiveUser(now.AddDate(0, -4, 0), promoRuleParams{}, rule, now)
|
||||
if eligible {
|
||||
t.Fatal("inactive promo should require positive inactive_months")
|
||||
}
|
||||
}
|
||||
|
||||
func TestLastSubscribeExpireAtIncludesUnlimitedSubscription(t *testing.T) {
|
||||
db, err := gorm.Open(mysql.New(mysql.Config{
|
||||
DSN: "gorm:password@tcp(localhost:9910)/gorm?charset=utf8&parseTime=True&loc=Local",
|
||||
SkipInitializeWithVersion: true,
|
||||
}), &gorm.Config{
|
||||
DryRun: true,
|
||||
DisableAutomaticPing: true,
|
||||
Logger: logger.Default.LogMode(logger.Silent),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("open gorm db: %v", err)
|
||||
}
|
||||
|
||||
stmt := lastSubscribeExpireQuery(db, 10).Take(nil).Statement
|
||||
sql := strings.ToLower(stmt.SQL.String())
|
||||
if strings.Contains(sql, "expire_time <>") || strings.Contains(sql, "expire_time !=") {
|
||||
t.Fatalf("last subscribe query should include unlimited subscription, sql: %s", stmt.SQL.String())
|
||||
}
|
||||
if !strings.Contains(sql, "case when expire_time = ? then 0 else 1 end asc") {
|
||||
t.Fatalf("last subscribe query should prioritize unlimited subscription, sql: %s", stmt.SQL.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscribePromoQueryMatchesQuantity(t *testing.T) {
|
||||
db, err := gorm.Open(mysql.New(mysql.Config{
|
||||
DSN: "gorm:password@tcp(localhost:9910)/gorm?charset=utf8&parseTime=True&loc=Local",
|
||||
SkipInitializeWithVersion: true,
|
||||
}), &gorm.Config{
|
||||
DryRun: true,
|
||||
DisableAutomaticPing: true,
|
||||
Logger: logger.Default.LogMode(logger.Silent),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("open gorm db: %v", err)
|
||||
}
|
||||
|
||||
stmt := subscribePromoQuery(db, 100, 3).Find(&[]subscribePromo{}).Statement
|
||||
sql := strings.ToLower(stmt.SQL.String())
|
||||
if !strings.Contains(sql, "subscribe_id = ?") {
|
||||
t.Fatalf("subscribe promo query should filter subscribe_id, sql: %s", stmt.SQL.String())
|
||||
}
|
||||
if !strings.Contains(sql, "quantity = ?") {
|
||||
t.Fatalf("subscribe promo query should filter quantity, sql: %s", stmt.SQL.String())
|
||||
}
|
||||
if got, want := len(stmt.Vars), 2; got != want {
|
||||
t.Fatalf("query vars len = %d, want %d, vars: %#v", got, want, stmt.Vars)
|
||||
}
|
||||
if stmt.Vars[0] != int64(100) || stmt.Vars[1] != int64(3) {
|
||||
t.Fatalf("query vars = %#v, want subscribe_id=100 quantity=3", stmt.Vars)
|
||||
}
|
||||
}
|
||||
@@ -154,12 +154,9 @@ func (l *PurchaseLogic) Purchase(req *types.PortalPurchaseRequest) (resp *types.
|
||||
}
|
||||
content, _ := tempOrder.Marshal()
|
||||
|
||||
// Persist activation context to DB so the worker can recover if Redis TTL expires.
|
||||
orderInfo.ActivationContext = string(content)
|
||||
|
||||
// Write to Redis as a hot cache (best-effort; non-fatal on failure).
|
||||
if _, redisErr := l.svcCtx.Redis.Set(l.ctx, fmt.Sprintf(constant.TempOrderCacheKey, orderInfo.OrderNo), string(content), CloseOrderTimeMinutes*time.Minute).Result(); redisErr != nil {
|
||||
l.Infow("[Purchase] Redis set error (non-fatal, DB fallback available)", logger.Field("error", redisErr.Error()), logger.Field("order_no", orderInfo.OrderNo))
|
||||
if _, err = l.svcCtx.Redis.Set(l.ctx, fmt.Sprintf(constant.TempOrderCacheKey, orderInfo.OrderNo), string(content), CloseOrderTimeMinutes*time.Minute).Result(); err != nil {
|
||||
l.Errorw("[Purchase] Redis set error", logger.Field("error", err.Error()), logger.Field("order_no", orderInfo.OrderNo))
|
||||
return err
|
||||
}
|
||||
l.Infow("[Purchase] Guest order", logger.Field("order_no", orderInfo.OrderNo), logger.Field("identifier", req.Identifier))
|
||||
|
||||
@@ -172,7 +169,7 @@ func (l *PurchaseLogic) Purchase(req *types.PortalPurchaseRequest) (resp *types.
|
||||
}
|
||||
}
|
||||
|
||||
// save guest order (activation_context is included)
|
||||
// save guest order
|
||||
if err = l.svcCtx.OrderModel.Insert(l.ctx, orderInfo, tx); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -151,48 +151,34 @@ func (l *RedeemCodeLogic) RedeemCode(req *types.RedeemCodeRequest) (resp *types.
|
||||
}
|
||||
|
||||
// 创建Order记录
|
||||
redemptionContext := struct {
|
||||
Type string `json:"type"`
|
||||
RedemptionCodeId int64 `json:"redemption_code_id"`
|
||||
UnitTime string `json:"unit_time"`
|
||||
Quantity int64 `json:"quantity"`
|
||||
}{
|
||||
Type: "redemption",
|
||||
RedemptionCodeId: redemptionCode.Id,
|
||||
UnitTime: redemptionCode.UnitTime,
|
||||
Quantity: redemptionCode.Quantity,
|
||||
}
|
||||
activationContextJSON, _ := json.Marshal(redemptionContext)
|
||||
|
||||
orderInfo := &order.Order{
|
||||
UserId: u.Id,
|
||||
OrderNo: tool.GenerateTradeNo(),
|
||||
Type: 5, // 兑换类型
|
||||
Quantity: redemptionCode.Quantity,
|
||||
Price: 0, // 兑换无价格
|
||||
Amount: 0, // 兑换无金额
|
||||
Discount: 0,
|
||||
GiftAmount: 0,
|
||||
Coupon: "",
|
||||
CouponDiscount: 0,
|
||||
PaymentId: 0,
|
||||
Method: "redemption",
|
||||
FeeAmount: 0,
|
||||
Commission: 0,
|
||||
Status: 2, // 直接设置为已支付
|
||||
SubscribeId: redemptionCode.SubscribePlan,
|
||||
IsNew: isNew,
|
||||
ActivationContext: string(activationContextJSON),
|
||||
UserId: u.Id,
|
||||
OrderNo: tool.GenerateTradeNo(),
|
||||
Type: 5, // 兑换类型
|
||||
Quantity: redemptionCode.Quantity,
|
||||
Price: 0, // 兑换无价格
|
||||
Amount: 0, // 兑换无金额
|
||||
Discount: 0,
|
||||
GiftAmount: 0,
|
||||
Coupon: "",
|
||||
CouponDiscount: 0,
|
||||
PaymentId: 0,
|
||||
Method: "redemption",
|
||||
FeeAmount: 0,
|
||||
Commission: 0,
|
||||
Status: 2, // 直接设置为已支付
|
||||
SubscribeId: redemptionCode.SubscribePlan,
|
||||
IsNew: isNew,
|
||||
}
|
||||
|
||||
// 保存Order到数据库(activation_context 同步写入,作为 Redis 的持久化兜底)
|
||||
// 保存Order到数据库
|
||||
err = l.svcCtx.OrderModel.Insert(l.ctx, orderInfo)
|
||||
if err != nil {
|
||||
l.Errorw("[RedeemCode] Create order failed", logger.Field("error", err.Error()))
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "create order failed")
|
||||
}
|
||||
|
||||
// 缓存兑换码信息到Redis(热缓存,供队列任务快速读取,非关键路径)
|
||||
// 缓存兑换码信息到Redis(供队列任务使用)
|
||||
cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo)
|
||||
cacheData := map[string]interface{}{
|
||||
"redemption_code_id": redemptionCode.Id,
|
||||
@@ -200,8 +186,16 @@ func (l *RedeemCodeLogic) RedeemCode(req *types.RedeemCodeRequest) (resp *types.
|
||||
"quantity": redemptionCode.Quantity,
|
||||
}
|
||||
jsonData, _ := json.Marshal(cacheData)
|
||||
if redisErr := l.svcCtx.Redis.Set(l.ctx, cacheKey, jsonData, 2*time.Hour).Err(); redisErr != nil {
|
||||
l.Infow("[RedeemCode] Cache redemption data failed (non-fatal, DB fallback available)", logger.Field("error", redisErr.Error()))
|
||||
err = l.svcCtx.Redis.Set(l.ctx, cacheKey, jsonData, 2*time.Hour).Err()
|
||||
if err != nil {
|
||||
l.Errorw("[RedeemCode] Cache redemption data failed", logger.Field("error", err.Error()))
|
||||
// 缓存失败,删除已创建的Order避免孤儿记录
|
||||
if delErr := l.svcCtx.OrderModel.Delete(l.ctx, orderInfo.Id); delErr != nil {
|
||||
l.Errorw("[RedeemCode] Delete order failed after cache error",
|
||||
logger.Field("order_id", orderInfo.Id),
|
||||
logger.Field("error", delErr.Error()))
|
||||
}
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.ERROR), "cache redemption data failed")
|
||||
}
|
||||
|
||||
// 触发队列任务
|
||||
|
||||
@@ -2,7 +2,10 @@ package user
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
logicCommon "github.com/perfect-panel/server/internal/logic/common"
|
||||
"github.com/perfect-panel/server/internal/model/log"
|
||||
"github.com/perfect-panel/server/internal/model/user"
|
||||
"github.com/perfect-panel/server/internal/svc"
|
||||
"github.com/perfect-panel/server/internal/types"
|
||||
@@ -35,38 +38,48 @@ func (l *CommissionWithdrawLogic) CommissionWithdraw(req *types.CommissionWithdr
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "Invalid Access")
|
||||
}
|
||||
|
||||
// Sum all pending (status=0) withdrawals to compute available balance.
|
||||
// Available = commission - pendingTotal; commission is only deducted on approval.
|
||||
var pendingTotal int64
|
||||
if err = l.svcCtx.DB.WithContext(l.ctx).
|
||||
Model(&user.Withdrawal{}).
|
||||
Where("user_id = ? AND status = 0", u.Id).
|
||||
Select("COALESCE(SUM(amount), 0)").
|
||||
Scan(&pendingTotal).Error; err != nil {
|
||||
l.Errorf("Failed to query pending withdrawals for user %d: %v", u.Id, err)
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "Failed to query pending withdrawals for user %d", u.Id)
|
||||
}
|
||||
|
||||
if u.Commission < req.Amount+pendingTotal {
|
||||
logger.Errorf("User %d insufficient available commission: total=%d pending=%d requested=%d",
|
||||
u.Id, u.Commission, pendingTotal, req.Amount)
|
||||
if u.Commission < req.Amount {
|
||||
logger.Errorf("User %d has insufficient commission balance: %.2f, requested: %.2f", u.Id, float64(u.Commission)/100, float64(req.Amount)/100)
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.UserCommissionNotEnough), "User %d has insufficient commission balance", u.Id)
|
||||
}
|
||||
|
||||
var w user.Withdrawal
|
||||
err = l.svcCtx.DB.WithContext(l.ctx).Transaction(func(tx *gorm.DB) error {
|
||||
w = user.Withdrawal{
|
||||
UserId: u.Id,
|
||||
Amount: req.Amount,
|
||||
Content: req.Content,
|
||||
Status: 0,
|
||||
Reason: "",
|
||||
}
|
||||
return tx.Create(&w).Error
|
||||
})
|
||||
tx := l.svcCtx.DB.WithContext(l.ctx).Begin()
|
||||
now := time.Now()
|
||||
|
||||
// Atomically deduct the requested amount so concurrent commission growth is preserved.
|
||||
if err = l.svcCtx.DB.WithContext(l.ctx).
|
||||
Model(&user.User{}).
|
||||
Where("id = ? AND commission >= ?", u.Id, req.Amount).
|
||||
UpdateColumn("commission", gorm.Expr("commission - ?", req.Amount)).Error; err != nil {
|
||||
tx.Rollback()
|
||||
l.Errorf("Failed to update user %d commission balance: %v", u.Id, err)
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "Failed to update user %d commission balance: %v", u.Id, err)
|
||||
}
|
||||
_ = l.svcCtx.UserModel.ClearUserCache(l.ctx, u)
|
||||
|
||||
// create withdrawal log
|
||||
if err = logicCommon.WriteCommissionLog(tx, u.Id, log.CommissionTypeConvertBalance, req.Amount, ""); err != nil {
|
||||
tx.Rollback()
|
||||
l.Errorf("Failed to create commission log for user %d: %v", u.Id, err)
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "Failed to create commission log for user %d: %v", u.Id, err)
|
||||
}
|
||||
|
||||
err = tx.Model(&user.Withdrawal{}).Create(&user.Withdrawal{
|
||||
UserId: u.Id,
|
||||
Amount: req.Amount,
|
||||
Content: req.Content,
|
||||
Status: 0,
|
||||
Reason: "",
|
||||
}).Error
|
||||
|
||||
if err != nil {
|
||||
l.Errorf("Failed to create withdrawal for user %d: %v", u.Id, err)
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "Failed to create withdrawal for user %d: %v", u.Id, err)
|
||||
tx.Rollback()
|
||||
l.Errorf("Failed to create withdrawal log for user %d: %v", u.Id, err)
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "Failed to create withdrawal log for user %d: %v", u.Id, err)
|
||||
}
|
||||
if err = tx.Commit().Error; err != nil {
|
||||
l.Errorf("Transaction commit failed for user %d withdrawal: %v", u.Id, err)
|
||||
return nil, errors.Wrapf(xerr.NewErrCode(xerr.ERROR), "Transaction commit failed for user %d withdrawal: %v", u.Id, err)
|
||||
}
|
||||
|
||||
return &types.WithdrawalLog{
|
||||
@@ -75,7 +88,7 @@ func (l *CommissionWithdrawLogic) CommissionWithdraw(req *types.CommissionWithdr
|
||||
Content: req.Content,
|
||||
Status: 0,
|
||||
Reason: "",
|
||||
CreatedAt: w.CreatedAt.UnixMilli(),
|
||||
UpdatedAt: w.UpdatedAt.UnixMilli(),
|
||||
CreatedAt: now.UnixMilli(),
|
||||
UpdatedAt: now.UnixMilli(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -46,7 +46,7 @@ func (l *DeviceWsConnectLogic) DeviceWsConnect(c *gin.Context) error {
|
||||
_, err := l.svcCtx.UserModel.FindOneDeviceByIdentifier(l.ctx, identifier)
|
||||
if err != nil && !sysErr.Is(err, gorm.ErrRecordNotFound) {
|
||||
l.Errorf("DeviceWsConnectLogic DeviceWsConnect FindOneDeviceByIdentifier err: %v", err)
|
||||
return errors.Wrap(xerr.NewErrCode(xerr.DatabaseQueryError), err.Error())
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), err.Error())
|
||||
}
|
||||
|
||||
value = l.ctx.Value(constant.CtxKeyUser)
|
||||
@@ -67,7 +67,7 @@ func (l *DeviceWsConnectLogic) DeviceWsConnect(c *gin.Context) error {
|
||||
err := l.svcCtx.UserModel.InsertDevice(l.ctx, &device)
|
||||
if err != nil {
|
||||
l.Errorf("DeviceWsConnectLogic DeviceWsConnect InsertDevice err: %v", err)
|
||||
return errors.Wrap(xerr.NewErrCode(xerr.DatabaseInsertError), err.Error())
|
||||
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), err.Error())
|
||||
}
|
||||
}
|
||||
//默认在线设备1
|
||||
|
||||
@@ -23,7 +23,6 @@ const (
|
||||
TypeSubscribeTraffic Type = 21 // Subscription traffic log
|
||||
TypeServerTraffic Type = 22 // Server traffic log
|
||||
TypeResetSubscribe Type = 23 // Reset subscription log
|
||||
TypeOrderRefund Type = 24 // Admin order refund log
|
||||
TypeLogin Type = 30 // Login log
|
||||
TypeRegister Type = 31 // Registration log
|
||||
TypeBalance Type = 32 // Balance log
|
||||
@@ -279,42 +278,6 @@ func (c *Commission) Unmarshal(data []byte) error {
|
||||
return json.Unmarshal(data, aux)
|
||||
}
|
||||
|
||||
type OrderRefund struct {
|
||||
OrderId int64 `json:"order_id"`
|
||||
OrderNo string `json:"order_no"`
|
||||
OperatorUserId int64 `json:"operator_user_id"`
|
||||
OperatorAuthIdentifier string `json:"operator_auth_identifier,omitempty"`
|
||||
TargetUserId int64 `json:"target_user_id"`
|
||||
UserSubscribeId int64 `json:"user_subscribe_id"`
|
||||
RefererUserId int64 `json:"referer_user_id,omitempty"`
|
||||
CommissionAmount int64 `json:"commission_amount"`
|
||||
Reason string `json:"reason,omitempty"`
|
||||
OrderStatusBefore uint8 `json:"order_status_before"`
|
||||
OrderStatusAfter uint8 `json:"order_status_after"`
|
||||
SubscribeStatusBefore uint8 `json:"subscribe_status_before"`
|
||||
SubscribeStatusAfter uint8 `json:"subscribe_status_after"`
|
||||
SubscribeExpireBefore int64 `json:"subscribe_expire_before"`
|
||||
SubscribeExpireAfter int64 `json:"subscribe_expire_after"`
|
||||
CommissionBefore int64 `json:"commission_before"`
|
||||
CommissionAfter int64 `json:"commission_after"`
|
||||
Timestamp int64 `json:"timestamp"`
|
||||
}
|
||||
|
||||
func (o *OrderRefund) Marshal() ([]byte, error) {
|
||||
type Alias OrderRefund
|
||||
return json.Marshal(&struct {
|
||||
*Alias
|
||||
}{
|
||||
Alias: (*Alias)(o),
|
||||
})
|
||||
}
|
||||
|
||||
func (o *OrderRefund) Unmarshal(data []byte) error {
|
||||
type Alias OrderRefund
|
||||
aux := (*Alias)(o)
|
||||
return json.Unmarshal(data, aux)
|
||||
}
|
||||
|
||||
// Gift represents a gift log entry.
|
||||
type Gift struct {
|
||||
Type uint16 `json:"type"`
|
||||
|
||||
@@ -160,8 +160,8 @@ func (m *customOrderModel) QueryMonthlyOrders(ctx context.Context, date time.Tim
|
||||
Where("status IN ? AND created_at BETWEEN ? AND ? AND method != ?", []int64{2, 5}, firstDay, lastDay, "balance").
|
||||
Select(
|
||||
"SUM(amount) as amount_total, " +
|
||||
"SUM(CASE WHEN type = 1 THEN amount ELSE 0 END) as new_order_amount, " +
|
||||
"SUM(CASE WHEN type = 2 THEN amount ELSE 0 END) as renewal_order_amount",
|
||||
"SUM(CASE WHEN is_new = 1 THEN amount ELSE 0 END) as new_order_amount, " +
|
||||
"SUM(CASE WHEN is_new = 0 THEN amount ELSE 0 END) as renewal_order_amount",
|
||||
).
|
||||
Scan(v).Error
|
||||
})
|
||||
@@ -177,8 +177,8 @@ func (m *customOrderModel) QueryDateOrders(ctx context.Context, date time.Time)
|
||||
Where("status IN ? AND DATE_FORMAT(created_at, '%Y-%m-%d') = ? AND method != ?", []int64{2, 5}, dateStr, "balance").
|
||||
Select(
|
||||
"SUM(amount) as amount_total, " +
|
||||
"SUM(CASE WHEN type = 1 THEN amount ELSE 0 END) as new_order_amount, " +
|
||||
"SUM(CASE WHEN type = 2 THEN amount ELSE 0 END) as renewal_order_amount",
|
||||
"SUM(CASE WHEN is_new = 1 THEN amount ELSE 0 END) as new_order_amount, " +
|
||||
"SUM(CASE WHEN is_new = 0 THEN amount ELSE 0 END) as renewal_order_amount",
|
||||
).
|
||||
Scan(v).Error
|
||||
})
|
||||
@@ -192,8 +192,8 @@ func (m *customOrderModel) QueryTotalOrders(ctx context.Context) (OrdersTotal, e
|
||||
return conn.Model(&Order{}).
|
||||
Select(`
|
||||
SUM(amount) AS amount_total,
|
||||
SUM(CASE WHEN type = 1 THEN amount ELSE 0 END) AS new_order_amount,
|
||||
SUM(CASE WHEN type = 2 THEN amount ELSE 0 END) AS renewal_order_amount
|
||||
SUM(CASE WHEN is_new = 1 THEN amount ELSE 0 END) AS new_order_amount,
|
||||
SUM(CASE WHEN is_new = 0 THEN amount ELSE 0 END) AS renewal_order_amount
|
||||
`).
|
||||
Where("status IN ? AND method != ?", []int64{2, 5}, "balance").
|
||||
Scan(&result).Error
|
||||
@@ -214,8 +214,8 @@ func (m *customOrderModel) QueryMonthlyUserCounts(ctx context.Context, date time
|
||||
err := m.QueryNoCacheCtx(ctx, nil, func(conn *gorm.DB, _ interface{}) error {
|
||||
return conn.Model(&Order{}).
|
||||
Select(`
|
||||
COUNT(DISTINCT CASE WHEN type = 1 THEN user_id END) AS new_users,
|
||||
COUNT(DISTINCT CASE WHEN type = 2 THEN user_id END) AS renewal_users
|
||||
COUNT(DISTINCT CASE WHEN is_new = 1 THEN user_id END) AS new_users,
|
||||
COUNT(DISTINCT CASE WHEN is_new = 0 THEN user_id END) AS renewal_users
|
||||
`).
|
||||
Where("status IN ? AND created_at >= ? AND created_at < ? AND method != ?",
|
||||
[]int64{2, 5}, firstDay, nextMonth, "balance").
|
||||
@@ -232,8 +232,8 @@ func (m *customOrderModel) QueryDateUserCounts(ctx context.Context, date time.Ti
|
||||
err := m.QueryNoCacheCtx(ctx, nil, func(conn *gorm.DB, _ interface{}) error {
|
||||
return conn.Model(&Order{}).
|
||||
Select(`
|
||||
COUNT(DISTINCT CASE WHEN type = 1 THEN user_id END) AS new_users,
|
||||
COUNT(DISTINCT CASE WHEN type = 2 THEN user_id END) AS renewal_users
|
||||
COUNT(DISTINCT CASE WHEN is_new = 1 THEN user_id END) AS new_users,
|
||||
COUNT(DISTINCT CASE WHEN is_new = 0 THEN user_id END) AS renewal_users
|
||||
`).
|
||||
Where("status IN ? AND DATE_FORMAT(created_at, '%Y-%m-%d') = ? AND method != ?",
|
||||
[]int64{2, 5}, dateStr, "balance").
|
||||
@@ -249,8 +249,8 @@ func (m *customOrderModel) QueryTotalUserCounts(ctx context.Context) (int64, int
|
||||
return conn.Model(&Order{}).
|
||||
Where("status IN ? AND method != ?", []int64{2, 5}, "balance").
|
||||
Select(`
|
||||
COUNT(DISTINCT CASE WHEN type = 1 THEN user_id END) AS new_users,
|
||||
COUNT(DISTINCT CASE WHEN type = 2 THEN user_id END) AS renewal_users
|
||||
COUNT(DISTINCT CASE WHEN is_new = 1 THEN user_id END) AS new_users,
|
||||
COUNT(DISTINCT CASE WHEN is_new = 0 THEN user_id END) AS renewal_users
|
||||
`).
|
||||
Scan(&counts).Error
|
||||
})
|
||||
@@ -282,8 +282,8 @@ func (m *customOrderModel) QueryDailyOrdersList(ctx context.Context, date time.T
|
||||
Select(`
|
||||
DATE_FORMAT(created_at, '%Y-%m-%d') AS date,
|
||||
SUM(amount) AS amount_total,
|
||||
SUM(CASE WHEN type = 1 THEN amount ELSE 0 END) AS new_order_amount,
|
||||
SUM(CASE WHEN type = 2 THEN amount ELSE 0 END) AS renewal_order_amount
|
||||
SUM(CASE WHEN is_new = 1 THEN amount ELSE 0 END) AS new_order_amount,
|
||||
SUM(CASE WHEN is_new = 0 THEN amount ELSE 0 END) AS renewal_order_amount
|
||||
`).
|
||||
Where("status IN ? AND created_at >= ? AND created_at < ? AND method != ?",
|
||||
[]int64{2, 5}, firstDay, nextDay, "balance").
|
||||
@@ -326,8 +326,8 @@ func (m *customOrderModel) QueryMonthlyOrdersList(ctx context.Context, date time
|
||||
Select(`
|
||||
DATE_FORMAT(created_at, '%Y-%m') AS date,
|
||||
SUM(amount) AS amount_total,
|
||||
SUM(CASE WHEN type = 1 THEN amount ELSE 0 END) AS new_order_amount,
|
||||
SUM(CASE WHEN type = 2 THEN amount ELSE 0 END) AS renewal_order_amount
|
||||
SUM(CASE WHEN is_new = 1 THEN amount ELSE 0 END) AS new_order_amount,
|
||||
SUM(CASE WHEN is_new = 0 THEN amount ELSE 0 END) AS renewal_order_amount
|
||||
`).
|
||||
Where("status IN ? AND created_at >= ? AND created_at < ? AND method != ?",
|
||||
[]int64{2, 5}, start, end, "balance").
|
||||
|
||||
@@ -24,9 +24,8 @@ type Order struct {
|
||||
Status uint8 `gorm:"type:tinyint(1);not null;default:1;comment:Order Status: 1: Pending, 2: Paid, 3:Close, 4: Failed, 5:Finished;"`
|
||||
SubscribeId int64 `gorm:"type:bigint;not null;default:0;comment:Subscribe Id"`
|
||||
SubscribeToken string `gorm:"type:varchar(255);default:null;comment:Renewal Subscribe Token"`
|
||||
AppAccountToken string `gorm:"type:varchar(36);default:null;comment:Apple IAP App Account Token (UUID)"`
|
||||
ActivationContext string `gorm:"type:text;default:null;comment:Activation context JSON (guest/redemption info for DB fallback)"`
|
||||
IsNew bool `gorm:"type:tinyint(1);not null;default:0;comment:Is New Order"`
|
||||
AppAccountToken string `gorm:"type:varchar(36);default:null;comment:Apple IAP App Account Token (UUID)"`
|
||||
IsNew bool `gorm:"type:tinyint(1);not null;default:0;comment:Is New Order"`
|
||||
CreatedAt time.Time `gorm:"<-:create;comment:Create Time"`
|
||||
UpdatedAt time.Time `gorm:"comment:Update Time"`
|
||||
}
|
||||
|
||||
@@ -378,13 +378,13 @@ func (m *customUserModel) QueryDailyUserStatisticsList(ctx context.Context, date
|
||||
// 子查询:统计每天的新用户订单数量
|
||||
newOrderSub := conn.Model(&order.Order{}).
|
||||
Select("DATE_FORMAT(created_at, '%Y-%m-%d') AS date, COUNT(DISTINCT user_id) AS new_order_users").
|
||||
Where("type = 1 AND created_at BETWEEN ? AND ? AND status IN ? AND method != ?", firstDay, date, []int64{2, 5}, "balance").
|
||||
Where("is_new = 1 AND created_at BETWEEN ? AND ? AND status IN ?", firstDay, date, []int64{2, 5}).
|
||||
Group("DATE_FORMAT(created_at, '%Y-%m-%d')")
|
||||
|
||||
// 子查询:统计每天的续费订单数量
|
||||
renewalOrderSub := conn.Model(&order.Order{}).
|
||||
Select("DATE_FORMAT(created_at, '%Y-%m-%d') AS date, COUNT(DISTINCT user_id) AS renewal_order_users").
|
||||
Where("type = 2 AND created_at BETWEEN ? AND ? AND status IN ? AND method != ?", firstDay, date, []int64{2, 5}, "balance").
|
||||
Where("is_new = 0 AND created_at BETWEEN ? AND ? AND status IN ?", firstDay, date, []int64{2, 5}).
|
||||
Group("DATE_FORMAT(created_at, '%Y-%m-%d')")
|
||||
|
||||
return conn.Model(&User{}).
|
||||
@@ -416,13 +416,13 @@ func (m *customUserModel) QueryMonthlyUserStatisticsList(ctx context.Context, da
|
||||
// 子查询:每月新订单用户数量
|
||||
newOrderSub := conn.Model(&order.Order{}).
|
||||
Select("DATE_FORMAT(created_at, '%Y-%m') AS date, COUNT(DISTINCT user_id) AS new_order_users").
|
||||
Where("type = 1 AND created_at >= ? AND status IN ? AND method != ?", sixMonthsAgo, []int64{2, 5}, "balance").
|
||||
Where("is_new = 1 AND created_at >= ? AND status IN ?", sixMonthsAgo, []int64{2, 5}).
|
||||
Group("DATE_FORMAT(created_at, '%Y-%m')")
|
||||
|
||||
// 子查询:每月续费订单用户数量
|
||||
renewalOrderSub := conn.Model(&order.Order{}).
|
||||
Select("DATE_FORMAT(created_at, '%Y-%m') AS date, COUNT(DISTINCT user_id) AS renewal_order_users").
|
||||
Where("type = 2 AND created_at >= ? AND status IN ? AND method != ?", sixMonthsAgo, []int64{2, 5}, "balance").
|
||||
Where("is_new = 0 AND created_at >= ? AND status IN ?", sixMonthsAgo, []int64{2, 5}).
|
||||
Group("DATE_FORMAT(created_at, '%Y-%m')")
|
||||
|
||||
return conn.Model(&User{}).
|
||||
|
||||
@@ -174,5 +174,5 @@ type Withdrawal struct {
|
||||
}
|
||||
|
||||
func (*Withdrawal) TableName() string {
|
||||
return "withdrawals"
|
||||
return "user_withdrawal"
|
||||
}
|
||||
|
||||
@@ -3,17 +3,10 @@
|
||||
|
||||
package types
|
||||
|
||||
import "encoding/json"
|
||||
|
||||
type ActivateOrderRequest struct {
|
||||
OrderNo string `json:"order_no" validate:"required"`
|
||||
}
|
||||
|
||||
type RefundOrderRequest struct {
|
||||
Id int64 `json:"id" validate:"required"`
|
||||
Reason string `json:"reason,omitempty" validate:"omitempty,max=500"`
|
||||
}
|
||||
|
||||
type Ads struct {
|
||||
Id int `json:"id"`
|
||||
Title string `json:"title"`
|
||||
@@ -833,17 +826,6 @@ type FilterCommissionLogResponse struct {
|
||||
List []CommissionLog `json:"list"`
|
||||
}
|
||||
|
||||
type FilterOrderRefundLogRequest struct {
|
||||
FilterLogParams
|
||||
OrderId int64 `form:"order_id,optional"`
|
||||
UserId int64 `form:"user_id,optional"`
|
||||
}
|
||||
|
||||
type FilterOrderRefundLogResponse struct {
|
||||
Total int64 `json:"total"`
|
||||
List []OrderRefundLog `json:"list"`
|
||||
}
|
||||
|
||||
type FilterEmailLogResponse struct {
|
||||
Total int64 `json:"total"`
|
||||
List []MessageLog `json:"list"`
|
||||
@@ -1870,7 +1852,6 @@ type Order struct {
|
||||
FeeAmount int64 `json:"fee_amount"`
|
||||
TradeNo string `json:"trade_no"`
|
||||
Status uint8 `json:"status"`
|
||||
StatusName string `json:"status_name,omitempty"`
|
||||
SubscribeId int64 `json:"subscribe_id"`
|
||||
CreatedAt int64 `json:"created_at"`
|
||||
UpdatedAt int64 `json:"updated_at"`
|
||||
@@ -1894,34 +1875,12 @@ type OrderDetail struct {
|
||||
FeeAmount int64 `json:"fee_amount"`
|
||||
TradeNo string `json:"trade_no"`
|
||||
Status uint8 `json:"status"`
|
||||
StatusName string `json:"status_name,omitempty"`
|
||||
SubscribeId int64 `json:"subscribe_id"`
|
||||
Subscribe Subscribe `json:"subscribe"`
|
||||
CreatedAt int64 `json:"created_at"`
|
||||
UpdatedAt int64 `json:"updated_at"`
|
||||
}
|
||||
|
||||
type OrderRefundLog struct {
|
||||
OrderId int64 `json:"order_id"`
|
||||
OrderNo string `json:"order_no"`
|
||||
OperatorUserId int64 `json:"operator_user_id"`
|
||||
OperatorAuthIdentifier string `json:"operator_auth_identifier,omitempty"`
|
||||
TargetUserId int64 `json:"target_user_id"`
|
||||
UserSubscribeId int64 `json:"user_subscribe_id"`
|
||||
RefererUserId int64 `json:"referer_user_id,omitempty"`
|
||||
CommissionAmount int64 `json:"commission_amount"`
|
||||
Reason string `json:"reason,omitempty"`
|
||||
OrderStatusBefore uint8 `json:"order_status_before"`
|
||||
OrderStatusAfter uint8 `json:"order_status_after"`
|
||||
SubscribeStatusBefore uint8 `json:"subscribe_status_before"`
|
||||
SubscribeStatusAfter uint8 `json:"subscribe_status_after"`
|
||||
SubscribeExpireBefore int64 `json:"subscribe_expire_before"`
|
||||
SubscribeExpireAfter int64 `json:"subscribe_expire_after"`
|
||||
CommissionBefore int64 `json:"commission_before"`
|
||||
CommissionAfter int64 `json:"commission_after"`
|
||||
Timestamp int64 `json:"timestamp"`
|
||||
}
|
||||
|
||||
type OrdersStatistics struct {
|
||||
Date string `json:"date,omitempty"`
|
||||
AmountTotal int64 `json:"amount_total"`
|
||||
@@ -3682,29 +3641,3 @@ type GetAdminUserInviteListResponse struct {
|
||||
Total int64 `json:"total"`
|
||||
List []AdminInvitedUser `json:"list"`
|
||||
}
|
||||
|
||||
type GetLogMessageRawRequest struct {
|
||||
Id int64 `form:"id" validate:"required"`
|
||||
}
|
||||
|
||||
type GetLogMessageRawResponse struct {
|
||||
Id int64 `json:"id"`
|
||||
Platform string `json:"platform"`
|
||||
AppVersion string `json:"app_version"`
|
||||
OsName string `json:"os_name"`
|
||||
OsVersion string `json:"os_version"`
|
||||
DeviceId string `json:"device_id"`
|
||||
UserId *int64 `json:"user_id"`
|
||||
SessionId string `json:"session_id"`
|
||||
Level uint8 `json:"level"`
|
||||
ErrorCode string `json:"error_code"`
|
||||
Message string `json:"message"`
|
||||
Stack string `json:"stack"`
|
||||
Context json.RawMessage `json:"context"`
|
||||
ClientIP string `json:"client_ip"`
|
||||
UserAgent string `json:"user_agent"`
|
||||
Locale string `json:"locale"`
|
||||
Digest string `json:"digest"`
|
||||
OccurredAt int64 `json:"occurred_at"`
|
||||
CreatedAt int64 `json:"created_at"`
|
||||
}
|
||||
|
||||
@@ -1,33 +0,0 @@
|
||||
# Docker Image Version Pins
|
||||
|
||||
This file records the infrastructure image versions pinned in `docker-compose.cloud.yml`.
|
||||
The versions below match the images observed on the test deployment on 2026-05-26.
|
||||
|
||||
| Service | Image | Running version source |
|
||||
| --- | --- | --- |
|
||||
| grafana | `grafana/grafana:13.0.1` | `grafana version 13.0.1` |
|
||||
| prometheus | `prom/prometheus:v3.11.3` | `prometheus, version 3.11.3` |
|
||||
| nginx-exporter | `nginx/nginx-prometheus-exporter:1.5.0` | image label `org.opencontainers.image.version=1.5.0` |
|
||||
| node-exporter | `prom/node-exporter:v1.11.1` | `node_exporter, version 1.11.1` |
|
||||
| cadvisor | `gcr.io/cadvisor/cadvisor:v0.55.1` | `cAdvisor version v0.55.1` |
|
||||
|
||||
The test deployment in `/root/bindbox/docker-compose.cloud.yml` also contains
|
||||
live-only exporter services that are not present in this repository's
|
||||
`docker-compose.cloud.yml`. They were pinned during staging validation:
|
||||
|
||||
| Test-only service | Image | Running version source |
|
||||
| --- | --- | --- |
|
||||
| mysql-exporter | `prom/mysqld-exporter:v0.19.0` | `mysqld_exporter, version 0.19.0` |
|
||||
| redis-exporter | `oliver006/redis_exporter:v1.82.0` | image label `org.opencontainers.image.version=v1.82.0` |
|
||||
|
||||
`ppanel-server` intentionally remains variable and requires `PPANEL_SERVER_TAG`
|
||||
from CI/CD so deployments use an immutable application image tag.
|
||||
|
||||
## Rollback
|
||||
|
||||
Restore the previous compose file from git and redeploy:
|
||||
|
||||
```sh
|
||||
git checkout HEAD~1 -- docker-compose.cloud.yml .env.example ops/docker-image-version-pins.md
|
||||
docker compose -f docker-compose.cloud.yml up -d
|
||||
```
|
||||
+5
-8
@@ -141,12 +141,9 @@ const (
|
||||
)
|
||||
|
||||
const (
|
||||
OrderNotExist uint32 = 61001
|
||||
PaymentMethodNotFound uint32 = 61002
|
||||
OrderStatusError uint32 = 61003
|
||||
InsufficientOfPeriod uint32 = 61004
|
||||
ExistAvailableTraffic uint32 = 61005
|
||||
OrderAlreadyRefunded uint32 = 61006
|
||||
OrderRefundNoSubscription uint32 = 61007
|
||||
OrderRefundCommissionMismatch uint32 = 61008
|
||||
OrderNotExist uint32 = 61001
|
||||
PaymentMethodNotFound uint32 = 61002
|
||||
OrderStatusError uint32 = 61003
|
||||
InsufficientOfPeriod uint32 = 61004
|
||||
ExistAvailableTraffic uint32 = 61005
|
||||
)
|
||||
|
||||
+4
-7
@@ -99,13 +99,10 @@ func init() {
|
||||
UseridNotMatch: "Userid not match",
|
||||
|
||||
// Order error
|
||||
OrderNotExist: "Order does not exist",
|
||||
PaymentMethodNotFound: "Payment method not found",
|
||||
OrderStatusError: "Order status error",
|
||||
InsufficientOfPeriod: "Insufficient number of period",
|
||||
OrderAlreadyRefunded: "Order already refunded",
|
||||
OrderRefundNoSubscription: "Refund target subscription not found",
|
||||
OrderRefundCommissionMismatch: "Refund commission source not found",
|
||||
OrderNotExist: "Order does not exist",
|
||||
PaymentMethodNotFound: "Payment method not found",
|
||||
OrderStatusError: "Order status error",
|
||||
InsufficientOfPeriod: "Insufficient number of period",
|
||||
|
||||
// Permission error
|
||||
PermissionDenied: "Permission denied",
|
||||
|
||||
@@ -48,7 +48,4 @@ 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
|
||||
mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx))
|
||||
}
|
||||
|
||||
@@ -7,7 +7,6 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/perfect-panel/server/internal/logic/admin/group"
|
||||
@@ -47,7 +46,7 @@ const (
|
||||
OrderStatusPaid = 2 // Order paid and ready for processing
|
||||
OrderStatusClose = 3 // Order closed/cancelled
|
||||
OrderStatusFailed = 4 // Order processing failed
|
||||
OrderStatusClaimed = 6 // Internal transient claim while a worker processes the order
|
||||
OrderStatusClaimed = 4 // Internal transient claim while a worker processes the order
|
||||
OrderStatusFinished = 5 // Order successfully completed
|
||||
)
|
||||
|
||||
@@ -93,12 +92,8 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
|
||||
|
||||
orderInfo, err := l.claimAndGetOrder(ctx, payload.OrderNo)
|
||||
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 // 返回错误触发 asynq 重试
|
||||
}
|
||||
logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过",
|
||||
logger.Field("order_no", payload.OrderNo))
|
||||
return nil
|
||||
@@ -121,12 +116,7 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
|
||||
)
|
||||
|
||||
if err = l.processOrderByType(ctx, orderInfo, payload.IAPExpireAt); err != nil {
|
||||
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()),
|
||||
)
|
||||
}
|
||||
l.releaseClaim(ctx, orderInfo.OrderNo)
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] 处理订单失败,将重试",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("order_type", orderInfo.Type),
|
||||
@@ -135,12 +125,7 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
|
||||
}
|
||||
|
||||
if err = l.reconcilePostOrderSubscriptions(ctx, orderInfo); err != nil {
|
||||
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()),
|
||||
)
|
||||
}
|
||||
l.releaseClaim(ctx, orderInfo.OrderNo)
|
||||
logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单订阅兜底合并失败,将重试",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("order_type", orderInfo.Type),
|
||||
@@ -191,15 +176,6 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// Detect stuck claimed order — return retryable error so asynq re-tries
|
||||
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),
|
||||
@@ -225,7 +201,7 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
|
||||
return &orderInfo, nil
|
||||
}
|
||||
|
||||
func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) error {
|
||||
func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) {
|
||||
if err := l.svc.DB.WithContext(ctx).
|
||||
Model(&order.Order{}).
|
||||
Where("order_no = ? AND status = ?", orderNo, OrderStatusClaimed).
|
||||
@@ -234,9 +210,7 @@ func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) e
|
||||
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
|
||||
@@ -589,27 +563,14 @@ func (l *ActivateOrderLogic) finalizeCouponAndOrder(ctx context.Context, orderIn
|
||||
}
|
||||
}
|
||||
|
||||
// UpdateOrderStatus uses WHERE status < target, which blocks claimed(6)→finished(5).
|
||||
// Use a direct update matching the exact claimed status, then update the full record
|
||||
// via the model layer to properly invalidate the cache.
|
||||
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 order cache regardless of whether the DB update succeeded
|
||||
orderInfo.Status = OrderStatusFinished
|
||||
if err := l.svc.OrderModel.Update(ctx, orderInfo); err != nil {
|
||||
logger.WithContext(ctx).Error("Update order cache after finalization failed",
|
||||
// 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()),
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
)
|
||||
}
|
||||
orderInfo.Status = OrderStatusFinished
|
||||
commonLogic.SubscriptionTraceInfo(logger.WithContext(ctx), commonLogic.SubscriptionTraceFlowOrder, "order_status_finished",
|
||||
"[SubscriptionFlow] order status updated to finished",
|
||||
commonLogic.OrderTraceFields(orderInfo)...,
|
||||
@@ -868,52 +829,26 @@ func (l *ActivateOrderLogic) createGuestUser(ctx context.Context, orderInfo *ord
|
||||
return userInfo, nil
|
||||
}
|
||||
|
||||
// getTempOrderInfo retrieves temporary order information from Redis cache with DB fallback.
|
||||
// getTempOrderInfo retrieves temporary order information from Redis cache
|
||||
func (l *ActivateOrderLogic) getTempOrderInfo(ctx context.Context, orderNo string) (*constant.TemporaryOrderInfo, error) {
|
||||
cacheKey := fmt.Sprintf(constant.TempOrderCacheKey, orderNo)
|
||||
data, err := l.svc.Redis.Get(ctx, cacheKey).Result()
|
||||
if err == nil {
|
||||
var tempOrder constant.TemporaryOrderInfo
|
||||
if unmarshalErr := tempOrder.Unmarshal([]byte(data)); unmarshalErr != nil {
|
||||
logger.WithContext(ctx).Error("Unmarshal temp order cache failed",
|
||||
logger.Field("error", unmarshalErr.Error()),
|
||||
logger.Field("cache_key", cacheKey),
|
||||
logger.Field("data", data),
|
||||
)
|
||||
return nil, unmarshalErr
|
||||
}
|
||||
return &tempOrder, nil
|
||||
}
|
||||
|
||||
// Redis miss — fall back to DB activation_context field.
|
||||
logger.WithContext(ctx).Infow("Redis cache miss for temp order, falling back to DB",
|
||||
logger.Field("cache_key", cacheKey),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
|
||||
orderInfo, dbErr := l.svc.OrderModel.FindOneByOrderNo(ctx, orderNo)
|
||||
if dbErr != nil {
|
||||
logger.WithContext(ctx).Error("DB fallback for temp order failed",
|
||||
logger.Field("order_no", orderNo),
|
||||
logger.Field("error", dbErr.Error()),
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("Get temp order cache failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("cache_key", cacheKey),
|
||||
)
|
||||
return nil, dbErr
|
||||
}
|
||||
|
||||
if orderInfo.ActivationContext == "" {
|
||||
logger.WithContext(ctx).Error("CRITICAL: activation_context missing in DB and Redis expired; cannot recover guest order",
|
||||
logger.Field("order_no", orderNo),
|
||||
)
|
||||
return nil, errors.New("activation context not found: Redis expired and no DB fallback available")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var tempOrder constant.TemporaryOrderInfo
|
||||
if unmarshalErr := tempOrder.Unmarshal([]byte(orderInfo.ActivationContext)); unmarshalErr != nil {
|
||||
logger.WithContext(ctx).Error("Unmarshal DB activation_context failed",
|
||||
logger.Field("order_no", orderNo),
|
||||
logger.Field("error", unmarshalErr.Error()),
|
||||
if err = tempOrder.Unmarshal([]byte(data)); err != nil {
|
||||
logger.WithContext(ctx).Error("Unmarshal temp order cache failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("cache_key", cacheKey),
|
||||
logger.Field("data", data),
|
||||
)
|
||||
return nil, unmarshalErr
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &tempOrder, nil
|
||||
@@ -1846,50 +1781,25 @@ func (l *ActivateOrderLogic) RedemptionActivate(ctx context.Context, orderInfo *
|
||||
return err
|
||||
}
|
||||
|
||||
// 3. 从Redis获取兑换码信息,Redis缺失时回查DB
|
||||
// 3. 从Redis获取兑换码信息
|
||||
cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo)
|
||||
data, err := l.svc.Redis.Get(ctx, cacheKey).Result()
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("Get redemption cache failed",
|
||||
logger.Field("error", err.Error()),
|
||||
logger.Field("cache_key", cacheKey),
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
var redemptionData struct {
|
||||
RedemptionCodeId int64 `json:"redemption_code_id"`
|
||||
UnitTime string `json:"unit_time"`
|
||||
Quantity int64 `json:"quantity"`
|
||||
}
|
||||
|
||||
cacheKey := fmt.Sprintf("redemption_order:%s", orderInfo.OrderNo)
|
||||
data, redisErr := l.svc.Redis.Get(ctx, cacheKey).Result()
|
||||
if redisErr == nil {
|
||||
if err = json.Unmarshal([]byte(data), &redemptionData); err != nil {
|
||||
logger.WithContext(ctx).Error("Unmarshal redemption cache failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
// Redis miss — fall back to DB activation_context field.
|
||||
logger.WithContext(ctx).Infow("Redis cache miss for redemption order, falling back to DB",
|
||||
logger.Field("cache_key", cacheKey),
|
||||
logger.Field("error", redisErr.Error()),
|
||||
)
|
||||
|
||||
if orderInfo.ActivationContext == "" {
|
||||
logger.WithContext(ctx).Error("CRITICAL: activation_context missing in DB and Redis expired; cannot recover redemption order",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
)
|
||||
return errors.New("activation context not found: Redis expired and no DB fallback available")
|
||||
}
|
||||
|
||||
var fullContext struct {
|
||||
Type string `json:"type"`
|
||||
RedemptionCodeId int64 `json:"redemption_code_id"`
|
||||
UnitTime string `json:"unit_time"`
|
||||
Quantity int64 `json:"quantity"`
|
||||
}
|
||||
if err = json.Unmarshal([]byte(orderInfo.ActivationContext), &fullContext); err != nil {
|
||||
logger.WithContext(ctx).Error("Unmarshal DB activation_context failed",
|
||||
logger.Field("order_no", orderInfo.OrderNo),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
return err
|
||||
}
|
||||
redemptionData.RedemptionCodeId = fullContext.RedemptionCodeId
|
||||
redemptionData.UnitTime = fullContext.UnitTime
|
||||
redemptionData.Quantity = fullContext.Quantity
|
||||
if err = json.Unmarshal([]byte(data), &redemptionData); err != nil {
|
||||
logger.WithContext(ctx).Error("Unmarshal redemption cache failed", logger.Field("error", err.Error()))
|
||||
return err
|
||||
}
|
||||
|
||||
// 4. 幂等性检查:查询是否已有兑换记录
|
||||
|
||||
@@ -1,104 +0,0 @@
|
||||
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"
|
||||
)
|
||||
|
||||
// StuckOrderRecoveryLogic scans orders stuck in claimed status and re-queues them for processing.
|
||||
type StuckOrderRecoveryLogic struct {
|
||||
svc *svc.ServiceContext
|
||||
}
|
||||
|
||||
func NewStuckOrderRecoveryLogic(svc *svc.ServiceContext) *StuckOrderRecoveryLogic {
|
||||
return &StuckOrderRecoveryLogic{svc: svc}
|
||||
}
|
||||
|
||||
// ProcessTask scans for orders stuck in claimed status for over 10 minutes,
|
||||
// resets them to paid, and re-enqueues the activate task so they are retried
|
||||
// independently of asynq's original retry counter (which may be exhausted).
|
||||
func (l *StuckOrderRecoveryLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error {
|
||||
cutoff := time.Now().Add(-10 * time.Minute)
|
||||
|
||||
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, recovering",
|
||||
logger.Field("count", len(stuckOrders)),
|
||||
logger.Field("order_nos", orderNos),
|
||||
)
|
||||
|
||||
for i := range stuckOrders {
|
||||
o := &stuckOrders[i]
|
||||
|
||||
result := l.svc.DB.WithContext(ctx).
|
||||
Model(&order.Order{}).
|
||||
Where("order_no = ? AND status = ?", o.OrderNo, OrderStatusClaimed).
|
||||
Update("status", OrderStatusPaid)
|
||||
if result.Error != nil {
|
||||
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to reset order status",
|
||||
logger.Field("order_no", o.OrderNo),
|
||||
logger.Field("error", result.Error.Error()),
|
||||
)
|
||||
continue
|
||||
}
|
||||
if result.RowsAffected == 0 {
|
||||
// Another process already handled this order
|
||||
continue
|
||||
}
|
||||
|
||||
// Invalidate order cache
|
||||
o.Status = OrderStatusPaid
|
||||
if err := l.svc.OrderModel.Update(ctx, o); err != nil {
|
||||
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to update order cache",
|
||||
logger.Field("order_no", o.OrderNo),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
}
|
||||
|
||||
// Re-enqueue activate task so the order gets processed regardless of asynq retry state
|
||||
payload, err := json.Marshal(queueTypes.ForthwithActivateOrderPayload{OrderNo: o.OrderNo})
|
||||
if err != nil {
|
||||
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to marshal task payload",
|
||||
logger.Field("order_no", o.OrderNo),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
continue
|
||||
}
|
||||
if _, err = l.svc.Queue.EnqueueContext(ctx, asynq.NewTask(queueTypes.ForthwithActivateOrder, payload)); err != nil {
|
||||
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to re-enqueue activate task",
|
||||
logger.Field("order_no", o.OrderNo),
|
||||
logger.Field("error", err.Error()),
|
||||
)
|
||||
} else {
|
||||
logger.WithContext(ctx).Info("[StuckOrderRecovery] Re-enqueued activate task",
|
||||
logger.Field("order_no", o.OrderNo),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -5,7 +5,6 @@ const (
|
||||
SchedulerTotalServerData = "scheduler:total:server"
|
||||
SchedulerResetTraffic = "scheduler:reset:traffic"
|
||||
SchedulerTrafficStat = "scheduler:traffic:stat"
|
||||
SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单
|
||||
SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账
|
||||
SchedulerStuckOrderRecovery = "scheduler:stuck:order:recovery" // 扫描并恢复超时 claimed 订单
|
||||
SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单
|
||||
SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账
|
||||
)
|
||||
|
||||
@@ -64,12 +64,6 @@ 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())
|
||||
}
|
||||
|
||||
@@ -1,5 +1,3 @@
|
||||
//go:build tools
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
|
||||
@@ -1,5 +1,3 @@
|
||||
//go:build tools
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
|
||||
Reference in New Issue
Block a user