All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m6s
111 lines
5.5 KiB
Go
111 lines
5.5 KiB
Go
package payment
|
|
|
|
import (
|
|
"context"
|
|
"time"
|
|
|
|
"github.com/bytedance/sonic"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
|
|
agentrecharge "github.com/break/junhong_cmp_fiber/internal/application/agentrecharge"
|
|
walletapp "github.com/break/junhong_cmp_fiber/internal/application/wallet"
|
|
"github.com/break/junhong_cmp_fiber/internal/infrastructure/messaging/outbox"
|
|
"github.com/break/junhong_cmp_fiber/internal/model"
|
|
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
|
"github.com/break/junhong_cmp_fiber/pkg/errors"
|
|
)
|
|
|
|
// AgentRechargePaymentConsumer 将已确认收款的代理在线充值幂等入账主钱包。
|
|
type AgentRechargePaymentConsumer struct {
|
|
db *gorm.DB
|
|
posting *walletapp.PostingService
|
|
}
|
|
|
|
// NewAgentRechargePaymentConsumer 创建代理在线充值入账消费者。
|
|
func NewAgentRechargePaymentConsumer(db *gorm.DB, posting *walletapp.PostingService) *AgentRechargePaymentConsumer {
|
|
return &AgentRechargePaymentConsumer{db: db, posting: posting}
|
|
}
|
|
|
|
// Consume 校验支付与充值权威事实后,在独立事务中完成唯一入账和充值终态。
|
|
func (c *AgentRechargePaymentConsumer) Consume(ctx context.Context, envelope outbox.DeliveryEnvelope) error {
|
|
if c == nil || c.db == nil || c.posting == nil {
|
|
return errors.New(errors.CodeInternalError, "代理在线充值入账消费者未配置")
|
|
}
|
|
if envelope.EventType != constants.OutboxEventTypeAgentRechargePaymentConfirmed ||
|
|
envelope.PayloadVersion != constants.AgentRechargePaymentConfirmedPayloadVersionV1 {
|
|
return errors.New(errors.CodeInvalidParam, "代理充值支付确认事件类型或版本不受支持")
|
|
}
|
|
var event agentrecharge.PaymentConfirmedEvent
|
|
if err := sonic.Unmarshal(envelope.Payload, &event); err != nil {
|
|
return errors.Wrap(errors.CodeInvalidParam, err, "代理充值支付确认事件载荷无法解析")
|
|
}
|
|
if event.EventID == "" || event.EventID != envelope.EventID || event.RechargeID == 0 || event.PaymentID == 0 ||
|
|
event.ShopID == 0 || event.WalletID == 0 || event.Amount <= 0 || event.ThirdPartyTradeNo == "" {
|
|
return errors.New(errors.CodeInvalidParam, "代理充值支付确认事件载荷不完整")
|
|
}
|
|
return c.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
recharge, _, err := lockCreditingFacts(ctx, tx, event)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if recharge.Status != constants.RechargeStatusPaid && recharge.Status != constants.RechargeStatusCompleted {
|
|
return errors.New(errors.CodeInvalidStatus, "代理在线充值当前状态不可入账")
|
|
}
|
|
if _, err := c.posting.PostInTx(ctx, tx, walletapp.PostingCommand{
|
|
ShopID: recharge.ShopID, WalletID: recharge.AgentWalletID, Amount: recharge.Amount,
|
|
ReferenceType: constants.ReferenceTypeTopup, ReferenceID: recharge.ID,
|
|
TransactionType: constants.AgentTransactionTypeRecharge,
|
|
UserID: recharge.UserID, Creator: recharge.UserID, Remark: "代理在线扫码充值",
|
|
RequestID: envelope.RequestID, CorrelationID: envelope.CorrelationID,
|
|
}); err != nil {
|
|
return err
|
|
}
|
|
if recharge.Status == constants.RechargeStatusCompleted {
|
|
return nil
|
|
}
|
|
completedAt := time.Now().UTC()
|
|
update := tx.WithContext(ctx).Model(&model.AgentRechargeRecord{}).
|
|
Where("id = ? AND status = ?", recharge.ID, constants.RechargeStatusPaid).
|
|
Updates(map[string]any{"status": constants.RechargeStatusCompleted, "completed_at": completedAt})
|
|
if update.Error != nil {
|
|
return errors.Wrap(errors.CodeDatabaseError, update.Error, "完成代理在线充值单失败")
|
|
}
|
|
if update.RowsAffected != 1 {
|
|
return errors.New(errors.CodeConflict, "代理在线充值状态已变化")
|
|
}
|
|
return nil
|
|
})
|
|
}
|
|
|
|
func lockCreditingFacts(ctx context.Context, tx *gorm.DB, event agentrecharge.PaymentConfirmedEvent) (*model.AgentRechargeRecord, *model.Payment, error) {
|
|
var recharge model.AgentRechargeRecord
|
|
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", event.RechargeID).First(&recharge).Error; err != nil {
|
|
if err == gorm.ErrRecordNotFound {
|
|
return nil, nil, errors.New(errors.CodeNotFound, "代理在线充值单不存在")
|
|
}
|
|
return nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "锁定代理在线充值单失败")
|
|
}
|
|
var payment model.Payment
|
|
if err := tx.WithContext(ctx).Where("id = ?", event.PaymentID).First(&payment).Error; err != nil {
|
|
if err == gorm.ErrRecordNotFound {
|
|
return nil, nil, errors.New(errors.CodeConflict, "代理在线充值支付单不存在")
|
|
}
|
|
return nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "查询代理在线充值支付单失败")
|
|
}
|
|
if recharge.ID != event.RechargeID || recharge.RechargeNo != event.RechargeNo || recharge.ShopID != event.ShopID ||
|
|
recharge.AgentWalletID != event.WalletID || recharge.UserID != event.UserID || recharge.Amount != event.Amount ||
|
|
recharge.PaymentTransactionID == nil || *recharge.PaymentTransactionID != event.ThirdPartyTradeNo ||
|
|
(recharge.PaymentMethod != constants.RechargeMethodWechat && recharge.PaymentMethod != constants.RechargeMethodAlipay) {
|
|
return nil, nil, errors.New(errors.CodeConflict, "代理充值事件与充值权威事实不一致")
|
|
}
|
|
if payment.OrderType != model.PaymentOrderTypeAgentRecharge || payment.OrderID != recharge.ID ||
|
|
payment.PaymentNo != event.PaymentNo || payment.PaymentMethod != event.PaymentMethod || payment.Amount != event.Amount ||
|
|
payment.Status != model.PaymentRecordStatusPaid || payment.ThirdPartyTradeNo != event.ThirdPartyTradeNo {
|
|
return nil, nil, errors.New(errors.CodeConflict, "代理充值事件与支付权威事实不一致")
|
|
}
|
|
return &recharge, &payment, nil
|
|
}
|
|
|
|
var _ outbox.EventConsumer = (*AgentRechargePaymentConsumer)(nil)
|