Files
junhong_cmp_fiber/internal/service/recharge_order/service.go
break 98c145fe70
Some checks failed
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Failing after 3m55s
实现支付商户池与微信授权配置
新增收款商户、商户池轮询、微信授权配置独立管理;三类新支付
(C端套餐购买、C端资产钱包充值、代理在线预存款充值)无条件
经商户池选择并冻结路由,无旧综合配置回退。merchant_id 为空
历史支付继续按 payment_config_id 双读。凭证版本化加载与
ID+版本缓存保证轮换一致性。删除商户池新支付创建开关及全部
引用。
2026-09-09 18:13:04 +08:00

505 lines
19 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package recharge_order
import (
"context"
"strconv"
"time"
merchantpayment "github.com/break/junhong_cmp_fiber/internal/application/merchantpayment"
"github.com/break/junhong_cmp_fiber/internal/infrastructure/audit"
"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/auditcontext"
"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
auditWriter *audit.Writer
}
// SetPaymentAudit 注入充值支付统一审计 Writer。
func (s *Service) SetPaymentAudit(writer *audit.Writer) {
s.auditWriter = writer
}
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
}
// 迟到首次成功必须允许 pending 与 failed 状态进入确认;冻结路由决定统计归属。
if payment.Status != model.PaymentRecordStatusPending && payment.Status != model.PaymentRecordStatusFailed {
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 {
// 乐观锁允许 pending 或 failed 首次进入 paid重复成功由 RowsAffected/唯一事实保护。
oldPaymentStatus := payment.Status
if oldPaymentStatus != model.PaymentRecordStatusPending && oldPaymentStatus != model.PaymentRecordStatusFailed {
return errors.New(errors.CodeInvalidStatus, "支付状态不允许处理")
}
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, "更新支付信息失败")
}
if err := merchantpayment.RecordFirstSuccess(ctx, tx, payment, now); err != nil {
return 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
}
if s.auditWriter == nil {
return errors.New(errors.CodeInvalidStatus, "充值支付统一审计接缝未配置")
}
afterPayment := *payment
afterPayment.Status = model.PaymentRecordStatusPaid
afterPayment.ThirdPartyTradeNo = transactionID
afterPayment.PaidAt = &now
paymentResource := audit.PaymentResource(&afterPayment, constants.AuditResourceRelationPrimary, constants.AuditResourceRolePaymentTarget,
map[string]any{"status": payment.Status, "third_party_trade_no": payment.ThirdPartyTradeNo, "paid_at": payment.PaidAt},
map[string]any{"status": afterPayment.Status, "third_party_trade_no": afterPayment.ThirdPartyTradeNo, "paid_at": afterPayment.PaidAt})
rechargeID := strconv.FormatUint(uint64(rechargeOrder.ID), 10)
rechargeResource := audit.ResourceInput{
Type: constants.AuditResourceRechargeOrder, ID: &rechargeID, Key: rechargeOrder.RechargeOrderNo, DisplayName: rechargeOrder.RechargeOrderNo,
Relation: constants.AuditResourceRelationAffected, Role: constants.AuditResourceRolePaymentBusinessOrder,
IdentitySnapshot: map[string]any{
"id": rechargeOrder.ID, "recharge_order_no": rechargeOrder.RechargeOrderNo, "user_id": rechargeOrder.UserID,
"asset_wallet_id": rechargeOrder.AssetWalletID, "resource_type": rechargeOrder.ResourceType,
"resource_id": rechargeOrder.ResourceID, "amount": rechargeOrder.Amount, "status": model.RechargeOrderStatusPaid,
},
BeforeData: map[string]any{"status": rechargeOrder.Status}, AfterData: map[string]any{"status": model.RechargeOrderStatusPaid},
SubjectVisibility: constants.AuditSubjectResult, SubjectSummary: "充值支付已到账",
}
transactionIDValue := strconv.FormatUint(uint64(transaction.ID), 10)
transactionResource := audit.ResourceInput{
Type: constants.AuditResourceAssetWalletTransaction, ID: &transactionIDValue, Key: transactionIDValue, DisplayName: "资产钱包流水 " + transactionIDValue,
Relation: constants.AuditResourceRelationAffected, Role: constants.AuditResourceRolePaymentWalletTransaction,
IdentitySnapshot: map[string]any{
"id": transaction.ID, "asset_wallet_id": transaction.AssetWalletID, "resource_type": transaction.ResourceType,
"resource_id": transaction.ResourceID, "transaction_type": transaction.TransactionType,
"reference_type": transaction.ReferenceType, "reference_no": transaction.ReferenceNo, "status": transaction.Status,
},
AfterData: map[string]any{"amount": transaction.Amount, "balance_before": transaction.BalanceBefore, "balance_after": transaction.BalanceAfter},
}
resources := []audit.ResourceInput{paymentResource, rechargeResource, transactionResource}
references, err := audit.AssetRechargeReferences(ctx, tx, rechargeOrder)
if err != nil {
return err
}
for i := range references {
if references[i].Type != constants.AuditResourceAssetWallet {
continue
}
references[i].Relation = constants.AuditResourceRelationAffected
references[i].BeforeData = map[string]any{"balance": balanceBefore}
references[i].AfterData = map[string]any{"balance": balanceBefore + rechargeOrder.Amount}
references[i].SubjectVisibility = constants.AuditSubjectResult
references[i].SubjectSummary = "充值支付已到账"
}
resources = append(resources, references...)
return s.auditWriter.Append(ctx, tx, audit.AppendInput{
ActionCode: constants.AuditActionPaymentConfirmed, Summary: "第三方支付确认资产充值已到账",
ScopeType: constants.AuditScopePlatform, Result: constants.AuditResultSuccess,
CorrelationID: payment.PaymentNo, Resources: resources,
})
})
if err != nil {
return err
}
linkedIDs := rechargeOrder.LinkedPackageIDs
if len(linkedIDs) > 0 {
linkage := auditcontext.From(ctx)
taskPayload := task.AutoPurchasePayload{
RechargeOrderID: rechargeOrder.ID, RequestID: linkage.RequestID,
CorrelationID: linkage.CorrelationID, ParentEventID: linkage.ParentEventID,
}
if err := s.queueClient.EnqueueTask(ctx, constants.TaskTypeAutoPurchaseAfterRecharge, taskPayload,
asynq.MaxRetry(3),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeAutoPurchaseAfterRecharge)),
); 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_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_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_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_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
}