重构充值订单模块:用 tb_recharge_order + tb_payment 替换 tb_asset_recharge_record
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m32s
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m32s
- 新增 RechargeOrder 和 Payment 模型及对应 Store - 新增 ClientRechargeOrderHandler 提供充值订单列表/详情接口 - 修改 client_wallet.go 使用新表读写充值数据 - 修改 callback/payment.go 将 CRCH 订单路由到 rechargeOrderService - 修改 client_order/service.go 的强充流程使用新表 - 修改 auto_purchase.go 从 tb_recharge_order 读取 linked_package_ids - 修改 order/service.go 的 WalletPay 使用 tb_payment 记录 - 修改 wechat_config_store.go 从 tb_recharge_order 统计待支付充值数 - 移除 AssetRechargeStore 和 AssetRechargeRecord 的注册引用 - 修复文档生成器缺失 ClientRechargeOrder handler - 状态枚举改为 0-based: Pending=0, Paid=1, Closed=2, Refunded=3
This commit is contained in:
@@ -51,7 +51,8 @@ type Service struct {
|
||||
assetService *asset.Service
|
||||
purchaseValidationService *purchase_validation.Service
|
||||
orderStore *postgres.OrderStore
|
||||
rechargeRecordStore *postgres.AssetRechargeStore
|
||||
rechargeOrderStore *postgres.RechargeOrderStore
|
||||
paymentStore *postgres.PaymentStore
|
||||
walletStore *postgres.AssetWalletStore
|
||||
personalDeviceStore *postgres.PersonalCustomerDeviceStore
|
||||
openIDStore *postgres.PersonalCustomerOpenIDStore
|
||||
@@ -73,7 +74,8 @@ func New(
|
||||
assetService *asset.Service,
|
||||
purchaseValidationService *purchase_validation.Service,
|
||||
orderStore *postgres.OrderStore,
|
||||
rechargeRecordStore *postgres.AssetRechargeStore,
|
||||
rechargeOrderStore *postgres.RechargeOrderStore,
|
||||
paymentStore *postgres.PaymentStore,
|
||||
walletStore *postgres.AssetWalletStore,
|
||||
personalDeviceStore *postgres.PersonalCustomerDeviceStore,
|
||||
openIDStore *postgres.PersonalCustomerOpenIDStore,
|
||||
@@ -93,7 +95,8 @@ func New(
|
||||
assetService: assetService,
|
||||
purchaseValidationService: purchaseValidationService,
|
||||
orderStore: orderStore,
|
||||
rechargeRecordStore: rechargeRecordStore,
|
||||
rechargeOrderStore: rechargeOrderStore,
|
||||
paymentStore: paymentStore,
|
||||
walletStore: walletStore,
|
||||
personalDeviceStore: personalDeviceStore,
|
||||
openIDStore: openIDStore,
|
||||
@@ -172,6 +175,27 @@ func (s *Service) CreateOrder(ctx context.Context, customerID uint, req *dto.Cli
|
||||
if lockAcquired {
|
||||
_ = s.redis.Del(skipCtx, lockKey).Err()
|
||||
}
|
||||
existingValue, err := s.redis.Get(skipCtx, redisKey).Result()
|
||||
if err == nil && strings.HasPrefix(existingValue, "CRCH") {
|
||||
existingOrder, queryErr := s.rechargeOrderStore.GetByRechargeOrderNo(skipCtx, existingValue)
|
||||
if queryErr == nil && existingOrder.Status == model.RechargeOrderStatusPending {
|
||||
return &dto.ClientCreateOrderResponse{
|
||||
OrderType: "recharge",
|
||||
Recharge: &dto.ClientRechargeInfo{
|
||||
RechargeID: existingOrder.ID,
|
||||
RechargeNo: existingOrder.RechargeOrderNo,
|
||||
Amount: existingOrder.Amount,
|
||||
Status: rechargeStatusToClientStatus(int(existingOrder.Status)),
|
||||
StatusName: clientRechargeStatusName(rechargeStatusToClientStatus(int(existingOrder.Status))),
|
||||
AutoPurchaseStatus: existingOrder.AutoPurchaseStatus,
|
||||
},
|
||||
Idempotent: true,
|
||||
}, nil
|
||||
}
|
||||
if queryErr == nil {
|
||||
_ = s.redis.Del(skipCtx, redisKey).Err()
|
||||
}
|
||||
}
|
||||
return nil, errors.New(errors.CodeTooManyRequests, "订单正在创建中,请勿重复提交")
|
||||
}
|
||||
|
||||
@@ -360,17 +384,26 @@ func (s *Service) createForceRechargeOrder(
|
||||
return nil, errors.Wrap(errors.CodeInternalError, err, "序列化关联套餐失败")
|
||||
}
|
||||
|
||||
carrierID := resourceID
|
||||
recharge := &model.AssetRechargeRecord{
|
||||
rechargeOrderNo := generateClientRechargeNo()
|
||||
var iotCardID *uint
|
||||
var deviceID *uint
|
||||
if validationResult.Card != nil {
|
||||
iotCardID = &validationResult.Card.ID
|
||||
}
|
||||
if validationResult.Device != nil {
|
||||
deviceID = &validationResult.Device.ID
|
||||
}
|
||||
|
||||
rechargeOrder := &model.RechargeOrder{
|
||||
RechargeOrderNo: rechargeOrderNo,
|
||||
UserID: customerID,
|
||||
AssetWalletID: wallet.ID,
|
||||
ResourceType: resourceType,
|
||||
ResourceID: resourceID,
|
||||
RechargeNo: generateClientRechargeNo(),
|
||||
IotCardID: iotCardID,
|
||||
DeviceID: deviceID,
|
||||
Amount: forceRecharge.ForceRechargeAmount,
|
||||
PaymentMethod: model.PaymentMethodWechat,
|
||||
PaymentConfigID: &activeConfig.ID,
|
||||
Status: 1,
|
||||
Status: model.RechargeOrderStatusPending,
|
||||
ShopIDTag: wallet.ShopIDTag,
|
||||
EnterpriseIDTag: wallet.EnterpriseIDTag,
|
||||
OperatorType: "personal_customer",
|
||||
@@ -378,33 +411,47 @@ func (s *Service) createForceRechargeOrder(
|
||||
LinkedPackageIDs: datatypes.JSON(linkedPackageIDs),
|
||||
LinkedOrderType: resolveOrderType(validationResult),
|
||||
LinkedCarrierType: assetInfo.AssetType,
|
||||
LinkedCarrierID: &carrierID,
|
||||
AutoPurchaseStatus: constants.AutoPurchaseStatusPending,
|
||||
LinkedCarrierID: &resourceID,
|
||||
AutoPurchaseStatus: model.AutoPurchaseStatusPending,
|
||||
}
|
||||
|
||||
if err := s.rechargeRecordStore.Create(ctx, recharge); err != nil {
|
||||
return nil, errors.Wrap(errors.CodeDatabaseError, err, "创建充值记录失败")
|
||||
if err := s.rechargeOrderStore.Create(ctx, rechargeOrder); err != nil {
|
||||
return nil, errors.Wrap(errors.CodeDatabaseError, err, "创建充值订单失败")
|
||||
}
|
||||
|
||||
paymentResult, err := paymentProvider.CreateJSAPIPayment(ctx, recharge.RechargeNo, "余额充值", openID, int(recharge.Amount))
|
||||
paymentNo := generateClientRechargeNo()
|
||||
payment := &model.Payment{
|
||||
PaymentNo: paymentNo,
|
||||
OrderID: rechargeOrder.ID,
|
||||
OrderType: model.PaymentOrderTypeRecharge,
|
||||
PaymentMethod: model.PaymentByWechat,
|
||||
Amount: forceRecharge.ForceRechargeAmount,
|
||||
Status: model.PaymentRecordStatusPending,
|
||||
PaymentConfigID: &activeConfig.ID,
|
||||
}
|
||||
|
||||
if err := s.paymentStore.Create(ctx, payment); err != nil {
|
||||
return nil, errors.Wrap(errors.CodeDatabaseError, err, "创建支付记录失败")
|
||||
}
|
||||
|
||||
paymentResult, err := paymentProvider.CreateJSAPIPayment(ctx, rechargeOrderNo, "余额充值", openID, int(rechargeOrder.Amount))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 微信支付创建成功后才标记幂等 key,确保支付失败时 defer 能正确清除 key
|
||||
s.markClientPurchaseCreated(ctx, redisKey, recharge.RechargeNo)
|
||||
s.markClientPurchaseCreated(ctx, redisKey, rechargeOrderNo)
|
||||
*created = true
|
||||
|
||||
clientStatus := rechargeStatusToClientStatus(recharge.Status)
|
||||
clientStatus := rechargeStatusToClientStatus(int(rechargeOrder.Status))
|
||||
return &dto.ClientCreateOrderResponse{
|
||||
OrderType: "recharge",
|
||||
Recharge: &dto.ClientRechargeInfo{
|
||||
RechargeID: recharge.ID,
|
||||
RechargeNo: recharge.RechargeNo,
|
||||
Amount: recharge.Amount,
|
||||
RechargeID: rechargeOrder.ID,
|
||||
RechargeNo: rechargeOrderNo,
|
||||
Amount: rechargeOrder.Amount,
|
||||
Status: clientStatus,
|
||||
StatusName: clientRechargeStatusName(clientStatus),
|
||||
AutoPurchaseStatus: recharge.AutoPurchaseStatus,
|
||||
AutoPurchaseStatus: rechargeOrder.AutoPurchaseStatus,
|
||||
},
|
||||
PayConfig: buildClientPayConfigFromResult(paymentResult),
|
||||
LinkedPackageInfo: buildLinkedPackageInfo(validationResult, forceRecharge),
|
||||
|
||||
@@ -39,6 +39,7 @@ type Service struct {
|
||||
orderItemStore *postgres.OrderItemStore
|
||||
agentWalletStore *postgres.AgentWalletStore
|
||||
assetWalletStore *postgres.AssetWalletStore
|
||||
paymentStore *postgres.PaymentStore
|
||||
purchaseValidationService *purchase_validation.Service
|
||||
shopPackageAllocationStore *postgres.ShopPackageAllocationStore
|
||||
shopSeriesAllocationStore *postgres.ShopSeriesAllocationStore
|
||||
@@ -65,6 +66,7 @@ func New(
|
||||
orderItemStore *postgres.OrderItemStore,
|
||||
agentWalletStore *postgres.AgentWalletStore,
|
||||
assetWalletStore *postgres.AssetWalletStore,
|
||||
paymentStore *postgres.PaymentStore,
|
||||
purchaseValidationService *purchase_validation.Service,
|
||||
shopPackageAllocationStore *postgres.ShopPackageAllocationStore,
|
||||
shopSeriesAllocationStore *postgres.ShopSeriesAllocationStore,
|
||||
@@ -89,6 +91,7 @@ func New(
|
||||
orderItemStore: orderItemStore,
|
||||
agentWalletStore: agentWalletStore,
|
||||
assetWalletStore: assetWalletStore,
|
||||
paymentStore: paymentStore,
|
||||
purchaseValidationService: purchaseValidationService,
|
||||
shopPackageAllocationStore: shopPackageAllocationStore,
|
||||
shopSeriesAllocationStore: shopSeriesAllocationStore,
|
||||
@@ -1188,6 +1191,18 @@ func (s *Service) unfreezeWalletForCancel(ctx context.Context, tx *gorm.DB, orde
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Service) createWalletPaymentRecord(tx *gorm.DB, order *model.Order, paymentMethod string) error {
|
||||
payment := &model.Payment{
|
||||
PaymentNo: order.OrderNo,
|
||||
OrderID: order.ID,
|
||||
OrderType: model.PaymentOrderTypePackage,
|
||||
PaymentMethod: paymentMethod,
|
||||
Amount: order.TotalAmount,
|
||||
Status: model.PaymentRecordStatusPaid,
|
||||
}
|
||||
return tx.Create(payment).Error
|
||||
}
|
||||
|
||||
func (s *Service) WalletPay(ctx context.Context, orderID uint, buyerType string, buyerID uint) error {
|
||||
order, err := s.orderStore.GetByID(ctx, orderID)
|
||||
if err != nil {
|
||||
@@ -1205,6 +1220,11 @@ func (s *Service) WalletPay(ctx context.Context, orderID uint, buyerType string,
|
||||
return errors.New(errors.CodeInvalidStatus, "代购订单无需支付")
|
||||
}
|
||||
|
||||
existingPayment, _ := s.paymentStore.GetByOrderIDAndStatus(ctx, orderID, model.PaymentRecordStatusPaid)
|
||||
if existingPayment != nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
var resourceType string
|
||||
var resourceID uint
|
||||
|
||||
@@ -1286,6 +1306,10 @@ func (s *Service) WalletPay(ctx context.Context, orderID uint, buyerType string,
|
||||
return errors.New(errors.CodeInsufficientBalance, "余额不足或并发冲突")
|
||||
}
|
||||
|
||||
if err := s.createWalletPaymentRecord(tx, order, model.PaymentByWallet); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return s.activatePackage(ctx, tx, order)
|
||||
})
|
||||
} else {
|
||||
@@ -1370,6 +1394,10 @@ func (s *Service) WalletPay(ctx context.Context, orderID uint, buyerType string,
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "创建扣款流水失败")
|
||||
}
|
||||
|
||||
if err := s.createWalletPaymentRecord(tx, order, model.PaymentByWallet); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return s.activatePackage(ctx, tx, order)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -3,18 +3,12 @@ package recharge
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"time"
|
||||
|
||||
"github.com/break/junhong_cmp_fiber/internal/model"
|
||||
"github.com/break/junhong_cmp_fiber/internal/model/dto"
|
||||
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
||||
"github.com/break/junhong_cmp_fiber/internal/task"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/errors"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/payment"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/queue"
|
||||
"github.com/hibiken/asynq"
|
||||
"go.uber.org/zap"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
@@ -40,7 +34,6 @@ type WechatConfigServiceInterface interface {
|
||||
// 负责资产钱包(IoT卡/设备)的充值订单创建、预检、支付回调处理等业务逻辑
|
||||
type Service struct {
|
||||
db *gorm.DB
|
||||
assetRechargeStore *postgres.AssetRechargeStore
|
||||
assetWalletStore *postgres.AssetWalletStore
|
||||
assetWalletTransactionStore *postgres.AssetWalletTransactionStore
|
||||
iotCardStore *postgres.IotCardStore
|
||||
@@ -50,14 +43,12 @@ type Service struct {
|
||||
commissionRecordStore *postgres.CommissionRecordStore
|
||||
wechatConfigService WechatConfigServiceInterface
|
||||
paymentLoader payment.PaymentConfigLoader
|
||||
queueClient *queue.Client // 任务队列客户端,用于触发自动购包任务
|
||||
logger *zap.Logger
|
||||
}
|
||||
|
||||
// New 创建充值服务实例
|
||||
func New(
|
||||
db *gorm.DB,
|
||||
assetRechargeStore *postgres.AssetRechargeStore,
|
||||
assetWalletStore *postgres.AssetWalletStore,
|
||||
assetWalletTransactionStore *postgres.AssetWalletTransactionStore,
|
||||
iotCardStore *postgres.IotCardStore,
|
||||
@@ -68,11 +59,9 @@ func New(
|
||||
wechatConfigService WechatConfigServiceInterface,
|
||||
paymentLoader payment.PaymentConfigLoader,
|
||||
logger *zap.Logger,
|
||||
queueClient *queue.Client,
|
||||
) *Service {
|
||||
return &Service{
|
||||
db: db,
|
||||
assetRechargeStore: assetRechargeStore,
|
||||
assetWalletStore: assetWalletStore,
|
||||
assetWalletTransactionStore: assetWalletTransactionStore,
|
||||
iotCardStore: iotCardStore,
|
||||
@@ -82,96 +71,10 @@ func New(
|
||||
commissionRecordStore: commissionRecordStore,
|
||||
wechatConfigService: wechatConfigService,
|
||||
paymentLoader: paymentLoader,
|
||||
queueClient: queueClient,
|
||||
logger: logger,
|
||||
}
|
||||
}
|
||||
|
||||
// Create 创建充值订单
|
||||
// 验证资源、金额范围、强充要求,生成订单号
|
||||
func (s *Service) Create(ctx context.Context, req *dto.CreateRechargeRequest, userID uint) (*dto.RechargeResponse, error) {
|
||||
// 1. 验证金额范围
|
||||
if req.Amount < constants.AssetRechargeMinAmount {
|
||||
return nil, errors.New(errors.CodeRechargeAmountInvalid, "充值金额不能低于1元")
|
||||
}
|
||||
if req.Amount > constants.AssetRechargeMaxAmount {
|
||||
return nil, errors.New(errors.CodeRechargeAmountInvalid, "充值金额不能超过100000元")
|
||||
}
|
||||
|
||||
// 2. 获取资源(卡或设备)钱包
|
||||
var wallet *model.AssetWallet
|
||||
var err error
|
||||
|
||||
if req.ResourceType == "iot_card" {
|
||||
wallet, err = s.assetWalletStore.GetByResourceTypeAndID(ctx, "iot_card", req.ResourceID)
|
||||
} else if req.ResourceType == "device" {
|
||||
wallet, err = s.assetWalletStore.GetByResourceTypeAndID(ctx, "device", req.ResourceID)
|
||||
} else {
|
||||
return nil, errors.New(errors.CodeInvalidParam, "无效的资源类型")
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil, errors.New(errors.CodeWalletNotFound, "钱包不存在")
|
||||
}
|
||||
return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询钱包失败")
|
||||
}
|
||||
|
||||
// 3. 验证强充要求
|
||||
forceReq, err := s.checkForceRechargeRequirement(ctx, req.ResourceType, req.ResourceID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if forceReq.NeedForceRecharge && req.Amount != forceReq.ForceRechargeAmount {
|
||||
return nil, errors.New(errors.CodeForceRechargeAmountMismatch,
|
||||
fmt.Sprintf("必须充值%d分才能满足强充要求", forceReq.ForceRechargeAmount))
|
||||
}
|
||||
|
||||
// 4. 生成充值订单号
|
||||
rechargeNo := s.generateRechargeNo()
|
||||
|
||||
// 5. 查询当前生效的支付配置
|
||||
var paymentConfigID *uint
|
||||
if req.PaymentMethod == "wechat" || req.PaymentMethod == "alipay" {
|
||||
activeConfig, err := s.wechatConfigService.GetActiveConfig(ctx)
|
||||
if err != nil {
|
||||
s.logger.Warn("查询生效支付配置失败", zap.Error(err))
|
||||
}
|
||||
if activeConfig != nil {
|
||||
paymentConfigID = &activeConfig.ID
|
||||
}
|
||||
}
|
||||
|
||||
// 6. 创建充值订单
|
||||
recharge := &model.AssetRechargeRecord{
|
||||
UserID: userID,
|
||||
AssetWalletID: wallet.ID,
|
||||
ResourceType: req.ResourceType,
|
||||
ResourceID: req.ResourceID,
|
||||
RechargeNo: rechargeNo,
|
||||
Amount: req.Amount,
|
||||
PaymentMethod: req.PaymentMethod,
|
||||
PaymentConfigID: paymentConfigID,
|
||||
Status: constants.RechargeStatusPending,
|
||||
ShopIDTag: wallet.ShopIDTag,
|
||||
EnterpriseIDTag: wallet.EnterpriseIDTag,
|
||||
}
|
||||
|
||||
if err := s.assetRechargeStore.Create(ctx, recharge); err != nil {
|
||||
return nil, errors.Wrap(errors.CodeDatabaseError, err, "创建充值订单失败")
|
||||
}
|
||||
|
||||
s.logger.Info("创建充值订单成功",
|
||||
zap.Uint("recharge_id", recharge.ID),
|
||||
zap.String("recharge_no", rechargeNo),
|
||||
zap.Int64("amount", req.Amount),
|
||||
zap.Uint("user_id", userID),
|
||||
)
|
||||
|
||||
return s.buildRechargeResponse(recharge), nil
|
||||
}
|
||||
|
||||
// GetRechargeCheck 充值预检
|
||||
// 返回强充要求、金额限制等信息
|
||||
func (s *Service) GetRechargeCheck(ctx context.Context, resourceType string, resourceID uint) (*ForceRechargeRequirement, error) {
|
||||
@@ -200,218 +103,6 @@ func (s *Service) GetRechargeCheck(ctx context.Context, resourceType string, res
|
||||
return s.checkForceRechargeRequirement(ctx, resourceType, resourceID)
|
||||
}
|
||||
|
||||
// GetByID 根据ID查询充值订单详情
|
||||
// 支持数据权限过滤
|
||||
func (s *Service) GetByID(ctx context.Context, id uint, userID uint) (*dto.RechargeResponse, error) {
|
||||
recharge, err := s.assetRechargeStore.GetByID(ctx, id)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil, errors.New(errors.CodeRechargeNotFound, "充值订单不存在")
|
||||
}
|
||||
return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询充值订单失败")
|
||||
}
|
||||
|
||||
// 数据权限检查:只能查看自己的充值订单
|
||||
if recharge.UserID != userID {
|
||||
return nil, errors.New(errors.CodeForbidden, "无权查看此充值订单")
|
||||
}
|
||||
|
||||
return s.buildRechargeResponse(recharge), nil
|
||||
}
|
||||
|
||||
// List 查询充值订单列表
|
||||
// 支持分页、筛选、数据权限
|
||||
func (s *Service) List(ctx context.Context, req *dto.RechargeListRequest, userID uint) (*dto.RechargeListResponse, error) {
|
||||
page := req.Page
|
||||
pageSize := req.PageSize
|
||||
if page == 0 {
|
||||
page = 1
|
||||
}
|
||||
if pageSize == 0 {
|
||||
pageSize = constants.DefaultPageSize
|
||||
}
|
||||
|
||||
params := &postgres.ListAssetRechargeParams{
|
||||
Page: page,
|
||||
PageSize: pageSize,
|
||||
UserID: &userID,
|
||||
}
|
||||
|
||||
if req.Status != nil {
|
||||
params.Status = req.Status
|
||||
}
|
||||
if req.WalletID != nil {
|
||||
walletID := *req.WalletID
|
||||
params.AssetWalletID = &walletID
|
||||
}
|
||||
if req.StartTime != nil {
|
||||
params.StartTime = req.StartTime
|
||||
}
|
||||
if req.EndTime != nil {
|
||||
params.EndTime = req.EndTime
|
||||
}
|
||||
|
||||
recharges, total, err := s.assetRechargeStore.List(ctx, params)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询充值订单列表失败")
|
||||
}
|
||||
|
||||
var list []*dto.RechargeResponse
|
||||
for _, r := range recharges {
|
||||
list = append(list, s.buildRechargeResponse(r))
|
||||
}
|
||||
|
||||
return &dto.RechargeListResponse{
|
||||
List: list,
|
||||
Total: total,
|
||||
Page: page,
|
||||
PageSize: pageSize,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// HandlePaymentCallback 支付回调处理
|
||||
// 支持幂等性检查、事务处理、更新余额、触发佣金
|
||||
// 验签由外层 Handler 负责;如需按 payment_config_id 加载配置,可通过 s.paymentLoader 获取
|
||||
func (s *Service) HandlePaymentCallback(ctx context.Context, rechargeNo string, paymentMethod string, paymentTransactionID string) error {
|
||||
// 1. 查询充值订单
|
||||
recharge, err := s.assetRechargeStore.GetByRechargeNo(ctx, rechargeNo)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return errors.New(errors.CodeRechargeNotFound, "充值订单不存在")
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询充值订单失败")
|
||||
}
|
||||
|
||||
// 2. 幂等性检查:已支付则直接返回成功
|
||||
if recharge.Status == constants.RechargeStatusPaid || recharge.Status == constants.RechargeStatusCompleted {
|
||||
s.logger.Info("充值订单已支付,跳过处理",
|
||||
zap.String("recharge_no", rechargeNo),
|
||||
zap.Int("status", recharge.Status),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
|
||||
// 3. 检查订单状态是否允许支付
|
||||
if recharge.Status != constants.RechargeStatusPending {
|
||||
return errors.New(errors.CodeInvalidStatus, "订单状态不允许支付")
|
||||
}
|
||||
|
||||
// 4. 获取钱包信息
|
||||
wallet, err := s.assetWalletStore.GetByID(ctx, recharge.AssetWalletID)
|
||||
if err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询钱包失败")
|
||||
}
|
||||
|
||||
// 5. 获取钱包对应的资源类型和ID(从充值记录中直接获取)
|
||||
resourceType := recharge.ResourceType
|
||||
resourceID := recharge.ResourceID
|
||||
|
||||
// 6. 事务处理:更新订单状态、增加余额、更新累计充值、触发佣金
|
||||
now := time.Now()
|
||||
err = s.db.Transaction(func(tx *gorm.DB) error {
|
||||
// 6.1 更新充值订单状态(带状态检查,使用事务内 tx 确保原子性)
|
||||
oldStatus := constants.RechargeStatusPending
|
||||
if err := s.assetRechargeStore.UpdateStatusWithOptimisticLockDB(ctx, tx, recharge.ID, &oldStatus, constants.RechargeStatusPaid, &now, nil); err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新充值订单状态失败")
|
||||
}
|
||||
|
||||
// 6.2 更新支付信息(使用事务内 tx)
|
||||
if err := s.assetRechargeStore.UpdatePaymentInfoWithDB(ctx, tx, recharge.ID, &paymentMethod, &paymentTransactionID); err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新支付信息失败")
|
||||
}
|
||||
|
||||
// 6.3 增加钱包余额(使用乐观锁)
|
||||
balanceBefore := wallet.Balance
|
||||
result := tx.Model(&model.AssetWallet{}).
|
||||
Where("id = ? AND version = ?", wallet.ID, wallet.Version).
|
||||
Updates(map[string]any{
|
||||
"balance": gorm.Expr("balance + ?", recharge.Amount),
|
||||
"version": gorm.Expr("version + 1"),
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新钱包余额失败")
|
||||
}
|
||||
if result.RowsAffected == 0 {
|
||||
return errors.New(errors.CodeInternalError, "钱包版本冲突,请重试")
|
||||
}
|
||||
|
||||
// 6.4 创建钱包交易记录(reference_no 存储充值单号,便于前端跳转)
|
||||
remark := "钱包充值"
|
||||
refType := "recharge"
|
||||
transaction := &model.AssetWalletTransaction{
|
||||
AssetWalletID: wallet.ID,
|
||||
ResourceType: resourceType,
|
||||
ResourceID: resourceID,
|
||||
UserID: recharge.UserID,
|
||||
TransactionType: "recharge",
|
||||
Amount: recharge.Amount,
|
||||
BalanceBefore: balanceBefore,
|
||||
BalanceAfter: balanceBefore + recharge.Amount,
|
||||
Status: 1,
|
||||
ReferenceType: &refType,
|
||||
ReferenceNo: &recharge.RechargeNo,
|
||||
Remark: &remark,
|
||||
ShopIDTag: wallet.ShopIDTag,
|
||||
EnterpriseIDTag: wallet.EnterpriseIDTag,
|
||||
}
|
||||
if err := tx.Create(transaction).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "创建钱包交易记录失败")
|
||||
}
|
||||
|
||||
// 6.5 更新累计充值
|
||||
if err := s.updateAccumulatedRechargeInTx(ctx, tx, resourceType, resourceID, recharge.Amount); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// 6.6 触发一次性佣金判断
|
||||
if err := s.triggerOneTimeCommissionIfNeededInTx(ctx, tx, resourceType, resourceID, recharge.Amount, recharge.UserID); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// 6.7 更新充值订单状态为已完成
|
||||
if err := tx.Model(&model.AssetRechargeRecord{}).
|
||||
Where("id = ?", recharge.ID).
|
||||
Update("status", constants.RechargeStatusCompleted).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新充值订单完成状态失败")
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// 充值成功后,如有关联套餐则触发自动购包(强充绑套餐功能)
|
||||
linkedIDs := recharge.LinkedPackageIDs
|
||||
if len(linkedIDs) > 0 && string(linkedIDs) != "[]" && string(linkedIDs) != "null" {
|
||||
if s.queueClient != nil {
|
||||
taskPayload := task.AutoPurchasePayload{RechargeRecordID: recharge.ID}
|
||||
if enqueueErr := s.queueClient.EnqueueTask(ctx, constants.TaskTypeAutoPurchaseAfterRecharge, taskPayload,
|
||||
asynq.Queue(constants.QueueDefault),
|
||||
asynq.MaxRetry(3),
|
||||
); enqueueErr != nil {
|
||||
s.logger.Error("自动购包任务入队失败",
|
||||
zap.Uint("recharge_id", recharge.ID),
|
||||
zap.Error(enqueueErr))
|
||||
// 不影响充值成功,仅记录日志
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
s.logger.Info("充值支付回调处理成功",
|
||||
zap.String("recharge_no", rechargeNo),
|
||||
zap.Int64("amount", recharge.Amount),
|
||||
zap.String("resource_type", resourceType),
|
||||
zap.Uint("resource_id", resourceID),
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// checkForceRechargeRequirement 检查强充要求
|
||||
// 根据资源类型和ID检查是否需要强制充值指定金额
|
||||
func (s *Service) checkForceRechargeRequirement(ctx context.Context, resourceType string, resourceID uint) (*ForceRechargeRequirement, error) {
|
||||
@@ -505,299 +196,3 @@ func (s *Service) checkForceRechargeRequirement(ctx context.Context, resourceTyp
|
||||
|
||||
// updateAccumulatedRechargeInTx 更新累计充值(事务内使用)
|
||||
// 同时更新旧的 accumulated_recharge 字段和新的 accumulated_recharge_by_series JSON 字段
|
||||
func (s *Service) updateAccumulatedRechargeInTx(ctx context.Context, tx *gorm.DB, resourceType string, resourceID uint, amount int64) error {
|
||||
if resourceType == "iot_card" {
|
||||
var card model.IotCard
|
||||
if err := tx.First(&card, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询IoT卡失败")
|
||||
}
|
||||
|
||||
if card.SeriesID != nil {
|
||||
if err := card.AddAccumulatedRechargeBySeries(*card.SeriesID, amount); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "更新卡按系列累计充值失败")
|
||||
}
|
||||
}
|
||||
|
||||
result := tx.Model(&model.IotCard{}).
|
||||
Where("id = ?", resourceID).
|
||||
Updates(map[string]any{
|
||||
"accumulated_recharge": gorm.Expr("accumulated_recharge + ?", amount),
|
||||
"accumulated_recharge_by_series": card.AccumulatedRechargeBySeriesJSON,
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新卡累计充值失败")
|
||||
}
|
||||
} else if resourceType == "device" {
|
||||
var device model.Device
|
||||
if err := tx.First(&device, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询设备失败")
|
||||
}
|
||||
|
||||
if device.SeriesID != nil {
|
||||
if err := device.AddAccumulatedRechargeBySeries(*device.SeriesID, amount); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "更新设备按系列累计充值失败")
|
||||
}
|
||||
}
|
||||
|
||||
result := tx.Model(&model.Device{}).
|
||||
Where("id = ?", resourceID).
|
||||
Updates(map[string]any{
|
||||
"accumulated_recharge": gorm.Expr("accumulated_recharge + ?", amount),
|
||||
"accumulated_recharge_by_series": device.AccumulatedRechargeBySeriesJSON,
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新设备累计充值失败")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// triggerOneTimeCommissionIfNeededInTx 触发一次性佣金(事务内使用)
|
||||
// 检查是否满足一次性佣金触发条件,满足则创建佣金记录并入账
|
||||
func (s *Service) triggerOneTimeCommissionIfNeededInTx(ctx context.Context, tx *gorm.DB, resourceType string, resourceID uint, rechargeAmount int64, userID uint) error {
|
||||
var seriesID *uint
|
||||
var accumulatedRecharge int64
|
||||
var firstCommissionPaid bool
|
||||
var shopID *uint
|
||||
|
||||
if resourceType == "iot_card" {
|
||||
var card model.IotCard
|
||||
if err := tx.First(&card, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询IoT卡失败")
|
||||
}
|
||||
seriesID = card.SeriesID
|
||||
shopID = card.ShopID
|
||||
if seriesID != nil {
|
||||
accumulatedRecharge = card.GetAccumulatedRechargeBySeries(*seriesID)
|
||||
firstCommissionPaid = card.IsFirstRechargeTriggeredBySeries(*seriesID)
|
||||
}
|
||||
} else if resourceType == "device" {
|
||||
var device model.Device
|
||||
if err := tx.First(&device, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询设备失败")
|
||||
}
|
||||
seriesID = device.SeriesID
|
||||
shopID = device.ShopID
|
||||
if seriesID != nil {
|
||||
accumulatedRecharge = device.GetAccumulatedRechargeBySeries(*seriesID)
|
||||
firstCommissionPaid = device.IsFirstRechargeTriggeredBySeries(*seriesID)
|
||||
}
|
||||
}
|
||||
|
||||
if seriesID == nil || firstCommissionPaid {
|
||||
return nil
|
||||
}
|
||||
|
||||
if shopID == nil {
|
||||
s.logger.Warn("资源未归属店铺,无法发放一次性佣金",
|
||||
zap.String("resource_type", resourceType),
|
||||
zap.Uint("resource_id", resourceID),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
|
||||
series, err := s.packageSeriesStore.GetByID(ctx, *seriesID)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询套餐系列失败")
|
||||
}
|
||||
|
||||
config, cfgErr := series.GetOneTimeCommissionConfig()
|
||||
if cfgErr != nil || config == nil || !config.Enable {
|
||||
return nil
|
||||
}
|
||||
|
||||
allocation, err := s.shopSeriesAllocationStore.GetByShopAndSeries(ctx, *shopID, *seriesID)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询系列分配失败")
|
||||
}
|
||||
|
||||
var rechargeAmountToCheck int64
|
||||
switch config.TriggerType {
|
||||
case model.OneTimeCommissionTriggerFirstRecharge:
|
||||
rechargeAmountToCheck = rechargeAmount
|
||||
default:
|
||||
rechargeAmountToCheck = accumulatedRecharge
|
||||
}
|
||||
|
||||
if rechargeAmountToCheck < config.Threshold {
|
||||
return nil
|
||||
}
|
||||
|
||||
commissionAmount := allocation.OneTimeCommissionAmount
|
||||
if commissionAmount <= 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
var commissionWallet model.AgentWallet
|
||||
if err := tx.Where("shop_id = ? AND wallet_type = ?", *shopID, constants.AgentWalletTypeCommission).
|
||||
First(&commissionWallet).Error; err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
s.logger.Warn("店铺佣金钱包不存在,跳过佣金发放",
|
||||
zap.Uint("shop_id", *shopID),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询店铺佣金钱包失败")
|
||||
}
|
||||
|
||||
var iotCardID, deviceID *uint
|
||||
if resourceType == "iot_card" {
|
||||
iotCardID = &resourceID
|
||||
} else {
|
||||
deviceID = &resourceID
|
||||
}
|
||||
|
||||
commissionRecord := &model.CommissionRecord{
|
||||
BaseModel: model.BaseModel{
|
||||
Creator: userID,
|
||||
Updater: userID,
|
||||
},
|
||||
ShopID: *shopID,
|
||||
IotCardID: iotCardID,
|
||||
DeviceID: deviceID,
|
||||
CommissionSource: model.CommissionSourceOneTime,
|
||||
Amount: commissionAmount,
|
||||
Status: constants.CommissionStatusReleased,
|
||||
Remark: "钱包充值触发一次性佣金",
|
||||
}
|
||||
|
||||
if err := tx.Create(commissionRecord).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "创建佣金记录失败")
|
||||
}
|
||||
|
||||
// 11. 佣金入账到店铺佣金钱包
|
||||
balanceBefore := commissionWallet.Balance
|
||||
result := tx.Model(&model.AgentWallet{}).
|
||||
Where("id = ? AND version = ?", commissionWallet.ID, commissionWallet.Version).
|
||||
Updates(map[string]any{
|
||||
"balance": gorm.Expr("balance + ?", commissionAmount),
|
||||
"version": gorm.Expr("version + 1"),
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新佣金钱包余额失败")
|
||||
}
|
||||
if result.RowsAffected == 0 {
|
||||
return errors.New(errors.CodeInternalError, "佣金钱包版本冲突,请重试")
|
||||
}
|
||||
|
||||
// 12. 更新佣金记录的入账后余额
|
||||
now := time.Now()
|
||||
if err := tx.Model(commissionRecord).Updates(map[string]any{
|
||||
"balance_after": balanceBefore + commissionAmount,
|
||||
"released_at": now,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新佣金记录失败")
|
||||
}
|
||||
|
||||
// 13. 创建佣金钱包交易记录
|
||||
remark := "一次性佣金入账(充值触发)"
|
||||
refType := "commission"
|
||||
commissionTransaction := &model.AgentWalletTransaction{
|
||||
AgentWalletID: commissionWallet.ID,
|
||||
ShopID: *shopID,
|
||||
UserID: userID,
|
||||
TransactionType: "commission",
|
||||
Amount: commissionAmount,
|
||||
BalanceBefore: balanceBefore,
|
||||
BalanceAfter: balanceBefore + commissionAmount,
|
||||
Status: 1,
|
||||
ReferenceType: &refType,
|
||||
ReferenceID: &commissionRecord.ID,
|
||||
Remark: &remark,
|
||||
ShopIDTag: *shopID,
|
||||
}
|
||||
if err := tx.Create(commissionTransaction).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "创建佣金钱包交易记录失败")
|
||||
}
|
||||
|
||||
// 14. 标记一次性佣金已发放
|
||||
if resourceType == "iot_card" {
|
||||
var card model.IotCard
|
||||
if err := tx.First(&card, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询IoT卡失败")
|
||||
}
|
||||
if err := card.SetFirstRechargeTriggeredBySeries(*seriesID, true); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "设置卡佣金发放状态失败")
|
||||
}
|
||||
if err := tx.Model(&model.IotCard{}).Where("id = ?", resourceID).
|
||||
Updates(map[string]any{
|
||||
"first_commission_paid": true,
|
||||
"first_recharge_triggered_by_series": card.FirstRechargeTriggeredBySeriesJSON,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新卡佣金发放状态失败")
|
||||
}
|
||||
} else {
|
||||
var device model.Device
|
||||
if err := tx.First(&device, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询设备失败")
|
||||
}
|
||||
if err := device.SetFirstRechargeTriggeredBySeries(*seriesID, true); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "设置设备佣金发放状态失败")
|
||||
}
|
||||
if err := tx.Model(&model.Device{}).Where("id = ?", resourceID).
|
||||
Updates(map[string]any{
|
||||
"first_commission_paid": true,
|
||||
"first_recharge_triggered_by_series": device.FirstRechargeTriggeredBySeriesJSON,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新设备佣金发放状态失败")
|
||||
}
|
||||
}
|
||||
|
||||
s.logger.Info("一次性佣金发放成功",
|
||||
zap.String("resource_type", resourceType),
|
||||
zap.Uint("resource_id", resourceID),
|
||||
zap.Uint("shop_id", *shopID),
|
||||
zap.Int64("commission_amount", commissionAmount),
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// generateRechargeNo 生成充值订单号
|
||||
// 格式: RCH + 14位时间戳 + 6位随机数
|
||||
func (s *Service) generateRechargeNo() string {
|
||||
now := time.Now()
|
||||
timestamp := now.Format("20060102150405")
|
||||
randomNum := rand.Intn(1000000)
|
||||
return fmt.Sprintf("RCH%s%06d", timestamp, randomNum)
|
||||
}
|
||||
|
||||
// buildRechargeResponse 构建充值订单响应
|
||||
func (s *Service) buildRechargeResponse(recharge *model.AssetRechargeRecord) *dto.RechargeResponse {
|
||||
statusText := ""
|
||||
switch recharge.Status {
|
||||
case constants.RechargeStatusPending:
|
||||
statusText = "待支付"
|
||||
case constants.RechargeStatusPaid:
|
||||
statusText = "已支付"
|
||||
case constants.RechargeStatusCompleted:
|
||||
statusText = "已完成"
|
||||
case constants.RechargeStatusClosed:
|
||||
statusText = "已关闭"
|
||||
case constants.RechargeStatusRefunded:
|
||||
statusText = "已退款"
|
||||
}
|
||||
|
||||
return &dto.RechargeResponse{
|
||||
ID: recharge.ID,
|
||||
RechargeNo: recharge.RechargeNo,
|
||||
UserID: recharge.UserID,
|
||||
WalletID: recharge.AssetWalletID,
|
||||
Amount: recharge.Amount,
|
||||
PaymentMethod: recharge.PaymentMethod,
|
||||
PaymentChannel: recharge.PaymentChannel,
|
||||
PaymentTransactionID: recharge.PaymentTransactionID,
|
||||
Status: recharge.Status,
|
||||
StatusText: statusText,
|
||||
PaidAt: recharge.PaidAt,
|
||||
CompletedAt: recharge.CompletedAt,
|
||||
CreatedAt: recharge.CreatedAt,
|
||||
UpdatedAt: recharge.UpdatedAt,
|
||||
}
|
||||
}
|
||||
|
||||
433
internal/service/recharge_order/service.go
Normal file
433
internal/service/recharge_order/service.go
Normal file
@@ -0,0 +1,433 @@
|
||||
package recharge_order
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/break/junhong_cmp_fiber/internal/model"
|
||||
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
||||
"github.com/break/junhong_cmp_fiber/internal/task"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/errors"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/queue"
|
||||
"github.com/hibiken/asynq"
|
||||
"go.uber.org/zap"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type Service struct {
|
||||
db *gorm.DB
|
||||
rechargeOrderStore *postgres.RechargeOrderStore
|
||||
paymentStore *postgres.PaymentStore
|
||||
assetWalletStore *postgres.AssetWalletStore
|
||||
assetWalletTransactionStore *postgres.AssetWalletTransactionStore
|
||||
iotCardStore *postgres.IotCardStore
|
||||
deviceStore *postgres.DeviceStore
|
||||
packageSeriesStore *postgres.PackageSeriesStore
|
||||
shopSeriesAllocationStore *postgres.ShopSeriesAllocationStore
|
||||
commissionRecordStore *postgres.CommissionRecordStore
|
||||
queueClient *queue.Client
|
||||
logger *zap.Logger
|
||||
}
|
||||
|
||||
func New(
|
||||
db *gorm.DB,
|
||||
rechargeOrderStore *postgres.RechargeOrderStore,
|
||||
paymentStore *postgres.PaymentStore,
|
||||
assetWalletStore *postgres.AssetWalletStore,
|
||||
assetWalletTransactionStore *postgres.AssetWalletTransactionStore,
|
||||
iotCardStore *postgres.IotCardStore,
|
||||
deviceStore *postgres.DeviceStore,
|
||||
packageSeriesStore *postgres.PackageSeriesStore,
|
||||
shopSeriesAllocationStore *postgres.ShopSeriesAllocationStore,
|
||||
commissionRecordStore *postgres.CommissionRecordStore,
|
||||
queueClient *queue.Client,
|
||||
logger *zap.Logger,
|
||||
) *Service {
|
||||
return &Service{
|
||||
db: db,
|
||||
rechargeOrderStore: rechargeOrderStore,
|
||||
paymentStore: paymentStore,
|
||||
assetWalletStore: assetWalletStore,
|
||||
assetWalletTransactionStore: assetWalletTransactionStore,
|
||||
iotCardStore: iotCardStore,
|
||||
deviceStore: deviceStore,
|
||||
packageSeriesStore: packageSeriesStore,
|
||||
shopSeriesAllocationStore: shopSeriesAllocationStore,
|
||||
commissionRecordStore: commissionRecordStore,
|
||||
queueClient: queueClient,
|
||||
logger: logger,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Service) HandlePaymentCallback(ctx context.Context, paymentNo string, paymentMethod string, transactionID string) error {
|
||||
payment, err := s.paymentStore.GetByPaymentNo(ctx, paymentNo)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return errors.New(errors.CodeNotFound, "支付记录不存在")
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询支付记录失败")
|
||||
}
|
||||
|
||||
if payment.Status == model.PaymentRecordStatusPaid {
|
||||
s.logger.Info("支付记录已处理,跳过",
|
||||
zap.String("payment_no", paymentNo),
|
||||
zap.Int("status", payment.Status),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
|
||||
if payment.Status != model.PaymentRecordStatusPending {
|
||||
return errors.New(errors.CodeInvalidStatus, "支付状态不允许处理")
|
||||
}
|
||||
|
||||
rechargeOrder, err := s.rechargeOrderStore.GetByID(ctx, payment.OrderID)
|
||||
if err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询充值订单失败")
|
||||
}
|
||||
|
||||
wallet, err := s.assetWalletStore.GetByID(ctx, rechargeOrder.AssetWalletID)
|
||||
if err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询钱包失败")
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
err = s.db.Transaction(func(tx *gorm.DB) error {
|
||||
oldPaymentStatus := model.PaymentRecordStatusPending
|
||||
if err := s.paymentStore.UpdateStatusWithOptimisticLockDB(ctx, tx, payment.ID, &oldPaymentStatus, model.PaymentRecordStatusPaid, &now); err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新支付状态失败")
|
||||
}
|
||||
|
||||
updates := map[string]interface{}{
|
||||
"third_party_trade_no": transactionID,
|
||||
}
|
||||
if err := tx.WithContext(ctx).Model(&model.Payment{}).Where("id = ?", payment.ID).Updates(updates).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新支付信息失败")
|
||||
}
|
||||
|
||||
oldRechargeStatus := model.RechargeOrderStatusPending
|
||||
if err := s.rechargeOrderStore.UpdateStatusWithOptimisticLockDB(ctx, tx, rechargeOrder.ID, &oldRechargeStatus, model.RechargeOrderStatusPaid, &now); err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新充值订单状态失败")
|
||||
}
|
||||
|
||||
balanceBefore := wallet.Balance
|
||||
result := tx.Model(&model.AssetWallet{}).
|
||||
Where("id = ? AND version = ?", wallet.ID, wallet.Version).
|
||||
Updates(map[string]any{
|
||||
"balance": gorm.Expr("balance + ?", rechargeOrder.Amount),
|
||||
"version": gorm.Expr("version + 1"),
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新钱包余额失败")
|
||||
}
|
||||
if result.RowsAffected == 0 {
|
||||
return errors.New(errors.CodeInternalError, "钱包版本冲突,请重试")
|
||||
}
|
||||
|
||||
refType := "recharge"
|
||||
transaction := &model.AssetWalletTransaction{
|
||||
AssetWalletID: wallet.ID,
|
||||
ResourceType: rechargeOrder.ResourceType,
|
||||
ResourceID: rechargeOrder.ResourceID,
|
||||
UserID: rechargeOrder.UserID,
|
||||
TransactionType: "recharge",
|
||||
Amount: rechargeOrder.Amount,
|
||||
BalanceBefore: balanceBefore,
|
||||
BalanceAfter: balanceBefore + rechargeOrder.Amount,
|
||||
Status: 1,
|
||||
ReferenceType: &refType,
|
||||
ReferenceNo: &paymentNo,
|
||||
ShopIDTag: wallet.ShopIDTag,
|
||||
EnterpriseIDTag: wallet.EnterpriseIDTag,
|
||||
}
|
||||
if err := tx.Create(transaction).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "创建钱包交易记录失败")
|
||||
}
|
||||
|
||||
if err := s.updateAccumulatedRechargeInTx(ctx, tx, rechargeOrder.ResourceType, rechargeOrder.ResourceID, rechargeOrder.Amount); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := s.triggerOneTimeCommissionIfNeededInTx(ctx, tx, rechargeOrder, rechargeOrder.Amount); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
linkedIDs := rechargeOrder.LinkedPackageIDs
|
||||
if len(linkedIDs) > 0 {
|
||||
taskPayload := task.AutoPurchasePayload{RechargeOrderID: rechargeOrder.ID}
|
||||
if err := s.queueClient.EnqueueTask(ctx, constants.TaskTypeAutoPurchaseAfterRecharge, taskPayload,
|
||||
asynq.Queue(constants.QueueDefault),
|
||||
asynq.MaxRetry(3),
|
||||
); err != nil {
|
||||
s.logger.Error("自动购包任务入队失败",
|
||||
zap.Uint("recharge_order_id", rechargeOrder.ID),
|
||||
zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
s.logger.Info("充值支付回调处理成功",
|
||||
zap.String("payment_no", paymentNo),
|
||||
zap.Int64("amount", rechargeOrder.Amount),
|
||||
zap.String("resource_type", rechargeOrder.ResourceType),
|
||||
zap.Uint("resource_id", rechargeOrder.ResourceID),
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Service) updateAccumulatedRechargeInTx(ctx context.Context, tx *gorm.DB, resourceType string, resourceID uint, amount int64) error {
|
||||
if resourceType == "iot_card" {
|
||||
var card model.IotCard
|
||||
if err := tx.First(&card, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询IoT卡失败")
|
||||
}
|
||||
|
||||
if card.SeriesID != nil {
|
||||
if err := card.AddAccumulatedRechargeBySeries(*card.SeriesID, amount); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "更新卡按系列累计充值失败")
|
||||
}
|
||||
}
|
||||
|
||||
result := tx.Model(&model.IotCard{}).
|
||||
Where("id = ?", resourceID).
|
||||
Updates(map[string]any{
|
||||
"accumulated_recharge": gorm.Expr("accumulated_recharge + ?", amount),
|
||||
"accumulated_recharge_by_series": card.AccumulatedRechargeBySeriesJSON,
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新卡累计充值失败")
|
||||
}
|
||||
} else if resourceType == "device" {
|
||||
var device model.Device
|
||||
if err := tx.First(&device, resourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询设备失败")
|
||||
}
|
||||
|
||||
if device.SeriesID != nil {
|
||||
if err := device.AddAccumulatedRechargeBySeries(*device.SeriesID, amount); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "更新设备按系列累计充值失败")
|
||||
}
|
||||
}
|
||||
|
||||
result := tx.Model(&model.Device{}).
|
||||
Where("id = ?", resourceID).
|
||||
Updates(map[string]any{
|
||||
"accumulated_recharge": gorm.Expr("accumulated_recharge + ?", amount),
|
||||
"accumulated_recharge_by_series": device.AccumulatedRechargeBySeriesJSON,
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新设备累计充值失败")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Service) triggerOneTimeCommissionIfNeededInTx(ctx context.Context, tx *gorm.DB, rechargeOrder *model.RechargeOrder, rechargeAmount int64) error {
|
||||
var seriesID *uint
|
||||
var accumulatedRecharge int64
|
||||
var firstCommissionPaid bool
|
||||
var shopID *uint
|
||||
|
||||
if rechargeOrder.ResourceType == "iot_card" {
|
||||
var card model.IotCard
|
||||
if err := tx.First(&card, rechargeOrder.ResourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询IoT卡失败")
|
||||
}
|
||||
seriesID = card.SeriesID
|
||||
shopID = card.ShopID
|
||||
if seriesID != nil {
|
||||
accumulatedRecharge = card.GetAccumulatedRechargeBySeries(*seriesID)
|
||||
firstCommissionPaid = card.IsFirstRechargeTriggeredBySeries(*seriesID)
|
||||
}
|
||||
} else if rechargeOrder.ResourceType == "device" {
|
||||
var device model.Device
|
||||
if err := tx.First(&device, rechargeOrder.ResourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询设备失败")
|
||||
}
|
||||
seriesID = device.SeriesID
|
||||
shopID = device.ShopID
|
||||
if seriesID != nil {
|
||||
accumulatedRecharge = device.GetAccumulatedRechargeBySeries(*seriesID)
|
||||
firstCommissionPaid = device.IsFirstRechargeTriggeredBySeries(*seriesID)
|
||||
}
|
||||
}
|
||||
|
||||
if seriesID == nil || firstCommissionPaid {
|
||||
return nil
|
||||
}
|
||||
|
||||
if shopID == nil {
|
||||
s.logger.Warn("资源未归属店铺,无法发放一次性佣金",
|
||||
zap.String("resource_type", rechargeOrder.ResourceType),
|
||||
zap.Uint("resource_id", rechargeOrder.ResourceID),
|
||||
)
|
||||
return nil
|
||||
}
|
||||
|
||||
series, err := s.packageSeriesStore.GetByID(ctx, *seriesID)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询套餐系列失败")
|
||||
}
|
||||
|
||||
config, cfgErr := series.GetOneTimeCommissionConfig()
|
||||
if cfgErr != nil || config == nil || !config.Enable {
|
||||
return nil
|
||||
}
|
||||
|
||||
allocation, err := s.shopSeriesAllocationStore.GetByShopAndSeries(ctx, *shopID, *seriesID)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询系列分配失败")
|
||||
}
|
||||
|
||||
var rechargeAmountToCheck int64
|
||||
switch config.TriggerType {
|
||||
case model.OneTimeCommissionTriggerFirstRecharge:
|
||||
rechargeAmountToCheck = rechargeAmount
|
||||
default:
|
||||
rechargeAmountToCheck = accumulatedRecharge
|
||||
}
|
||||
|
||||
if rechargeAmountToCheck < config.Threshold {
|
||||
return nil
|
||||
}
|
||||
|
||||
commissionAmount := allocation.OneTimeCommissionAmount
|
||||
if commissionAmount <= 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
var commissionWallet model.AgentWallet
|
||||
if err := tx.Where("shop_id = ? AND wallet_type = ?", *shopID, constants.AgentWalletTypeCommission).
|
||||
First(&commissionWallet).Error; err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
s.logger.Warn("店铺佣金钱包不存在,跳过佣金发放", zap.Uint("shop_id", *shopID))
|
||||
return nil
|
||||
}
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询店铺佣金钱包失败")
|
||||
}
|
||||
|
||||
var iotCardID, deviceID *uint
|
||||
if rechargeOrder.ResourceType == "iot_card" {
|
||||
iotCardID = &rechargeOrder.ResourceID
|
||||
} else {
|
||||
deviceID = &rechargeOrder.ResourceID
|
||||
}
|
||||
|
||||
commissionRecord := &model.CommissionRecord{
|
||||
BaseModel: model.BaseModel{
|
||||
Creator: rechargeOrder.UserID,
|
||||
Updater: rechargeOrder.UserID,
|
||||
},
|
||||
ShopID: *shopID,
|
||||
IotCardID: iotCardID,
|
||||
DeviceID: deviceID,
|
||||
CommissionSource: model.CommissionSourceOneTime,
|
||||
Amount: commissionAmount,
|
||||
Status: constants.CommissionStatusReleased,
|
||||
Remark: "钱包充值触发一次性佣金",
|
||||
}
|
||||
|
||||
if err := tx.Create(commissionRecord).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "创建佣金记录失败")
|
||||
}
|
||||
|
||||
balanceBefore := commissionWallet.Balance
|
||||
result := tx.Model(&model.AgentWallet{}).
|
||||
Where("id = ? AND version = ?", commissionWallet.ID, commissionWallet.Version).
|
||||
Updates(map[string]any{
|
||||
"balance": gorm.Expr("balance + ?", commissionAmount),
|
||||
"version": gorm.Expr("version + 1"),
|
||||
})
|
||||
if result.Error != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新佣金钱包余额失败")
|
||||
}
|
||||
if result.RowsAffected == 0 {
|
||||
return errors.New(errors.CodeInternalError, "佣金钱包版本冲突,请重试")
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
if err := tx.Model(commissionRecord).Updates(map[string]any{
|
||||
"balance_after": balanceBefore + commissionAmount,
|
||||
"released_at": now,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新佣金记录失败")
|
||||
}
|
||||
|
||||
refType := "commission"
|
||||
commissionTransaction := &model.AgentWalletTransaction{
|
||||
AgentWalletID: commissionWallet.ID,
|
||||
ShopID: *shopID,
|
||||
UserID: rechargeOrder.UserID,
|
||||
TransactionType: "commission",
|
||||
Amount: commissionAmount,
|
||||
BalanceBefore: balanceBefore,
|
||||
BalanceAfter: balanceBefore + commissionAmount,
|
||||
Status: 1,
|
||||
ReferenceType: &refType,
|
||||
ReferenceID: &commissionRecord.ID,
|
||||
ShopIDTag: *shopID,
|
||||
}
|
||||
if err := tx.Create(commissionTransaction).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "创建佣金钱包交易记录失败")
|
||||
}
|
||||
|
||||
if rechargeOrder.ResourceType == "iot_card" {
|
||||
var card model.IotCard
|
||||
if err := tx.First(&card, rechargeOrder.ResourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询IoT卡失败")
|
||||
}
|
||||
if err := card.SetFirstRechargeTriggeredBySeries(*seriesID, true); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "设置卡佣金发放状态失败")
|
||||
}
|
||||
if err := tx.Model(&model.IotCard{}).Where("id = ?", rechargeOrder.ResourceID).
|
||||
Updates(map[string]any{
|
||||
"first_commission_paid": true,
|
||||
"first_recharge_triggered_by_series": card.FirstRechargeTriggeredBySeriesJSON,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新卡佣金发放状态失败")
|
||||
}
|
||||
} else {
|
||||
var device model.Device
|
||||
if err := tx.First(&device, rechargeOrder.ResourceID).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "查询设备失败")
|
||||
}
|
||||
if err := device.SetFirstRechargeTriggeredBySeries(*seriesID, true); err != nil {
|
||||
return errors.Wrap(errors.CodeInternalError, err, "设置设备佣金发放状态失败")
|
||||
}
|
||||
if err := tx.Model(&model.Device{}).Where("id = ?", rechargeOrder.ResourceID).
|
||||
Updates(map[string]any{
|
||||
"first_commission_paid": true,
|
||||
"first_recharge_triggered_by_series": device.FirstRechargeTriggeredBySeriesJSON,
|
||||
}).Error; err != nil {
|
||||
return errors.Wrap(errors.CodeDatabaseError, err, "更新设备佣金发放状态失败")
|
||||
}
|
||||
}
|
||||
|
||||
s.logger.Info("一次性佣金发放成功",
|
||||
zap.String("resource_type", rechargeOrder.ResourceType),
|
||||
zap.Uint("resource_id", rechargeOrder.ResourceID),
|
||||
zap.Uint("shop_id", *shopID),
|
||||
zap.Int64("commission_amount", commissionAmount),
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -33,7 +33,7 @@ type AuditServiceInterface interface {
|
||||
type Service struct {
|
||||
store *postgres.WechatConfigStore
|
||||
orderStore *postgres.OrderStore
|
||||
assetRechargeStore *postgres.AssetRechargeStore
|
||||
rechargeOrderStore *postgres.RechargeOrderStore
|
||||
agentRechargeStore *postgres.AgentRechargeStore
|
||||
auditService AuditServiceInterface
|
||||
redis *redis.Client
|
||||
@@ -44,7 +44,7 @@ type Service struct {
|
||||
func New(
|
||||
store *postgres.WechatConfigStore,
|
||||
orderStore *postgres.OrderStore,
|
||||
assetRechargeStore *postgres.AssetRechargeStore,
|
||||
rechargeOrderStore *postgres.RechargeOrderStore,
|
||||
agentRechargeStore *postgres.AgentRechargeStore,
|
||||
auditService AuditServiceInterface,
|
||||
rdb *redis.Client,
|
||||
@@ -53,7 +53,7 @@ func New(
|
||||
return &Service{
|
||||
store: store,
|
||||
orderStore: orderStore,
|
||||
assetRechargeStore: assetRechargeStore,
|
||||
rechargeOrderStore: rechargeOrderStore,
|
||||
agentRechargeStore: agentRechargeStore,
|
||||
auditService: auditService,
|
||||
redis: rdb,
|
||||
@@ -521,7 +521,7 @@ func (s *Service) resolvePaymentConfigID(ctx context.Context, orderNo string) (*
|
||||
return order.PaymentConfigID, nil
|
||||
|
||||
case len(orderNo) >= 4 && orderNo[:4] == constants.AssetRechargeOrderPrefix:
|
||||
record, err := s.assetRechargeStore.GetByRechargeNo(ctx, orderNo)
|
||||
record, err := s.rechargeOrderStore.GetByRechargeOrderNo(ctx, orderNo)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user