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.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 }