2 Commits

Author SHA1 Message Date
586a1cccd5 修复佣金回扣问题
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m25s
2026-08-14 09:08:43 +08:00
b5285877bf 归档 2026-08-13 17:45:49 +08:00
14 changed files with 200 additions and 8 deletions

View File

@@ -117,7 +117,7 @@ func Recover(ctx context.Context, db *gorm.DB, repository *outbox.Repository, li
logger.Warn("扫描待计算订单失败", zap.Error(err)) logger.Warn("扫描待计算订单失败", zap.Error(err))
} else { } else {
for _, order := range orders { for _, order := range orders {
recoverOne(ctx, db, repository, EventCommissionCalculate, order.ID, order.ID, &orderStats, logger) recoverOne(ctx, db, repository, EventCommissionCalculate, order.ID, order.ID, &orderStats, logger, false)
} }
} }
var refunds []model.RefundRequest var refunds []model.RefundRequest
@@ -126,10 +126,10 @@ func Recover(ctx context.Context, db *gorm.DB, repository *outbox.Repository, li
} else { } else {
for _, refund := range refunds { for _, refund := range refunds {
if !refund.CommissionDeducted { if !refund.CommissionDeducted {
recoverOne(ctx, db, repository, EventRefundCommissionDeduct, refund.ID, refund.OrderID, &refundStats, logger) recoverOne(ctx, db, repository, EventRefundCommissionDeduct, refund.ID, refund.OrderID, &refundStats, logger, true)
} }
if !refund.AssetReset { if !refund.AssetReset {
recoverOne(ctx, db, repository, EventRefundAssetProcess, refund.ID, refund.OrderID, &refundStats, logger) recoverOne(ctx, db, repository, EventRefundAssetProcess, refund.ID, refund.OrderID, &refundStats, logger, true)
} }
} }
} }
@@ -138,16 +138,21 @@ func Recover(ctx context.Context, db *gorm.DB, repository *outbox.Repository, li
zap.Int("退款已补发", refundStats.Resent), zap.Int("退款无需补发", refundStats.Unchanged), zap.Int("退款失败", refundStats.Failed)) zap.Int("退款已补发", refundStats.Resent), zap.Int("退款无需补发", refundStats.Unchanged), zap.Int("退款失败", refundStats.Failed))
} }
func recoverOne(ctx context.Context, db *gorm.DB, repository *outbox.Repository, eventType string, aggregateID, orderID uint, stats *RecoveryStats, logger *zap.Logger) { func recoverOne(ctx context.Context, db *gorm.DB, repository *outbox.Repository, eventType string, aggregateID, orderID uint, stats *RecoveryStats, logger *zap.Logger, retryDelivered bool) {
eventID := outboxid.Stable(eventType+":", strconv.FormatUint(uint64(aggregateID), 10)) eventID := outboxid.Stable(eventType+":", strconv.FormatUint(uint64(aggregateID), 10))
var event model.OutboxEvent var event model.OutboxEvent
err := db.WithContext(ctx).Where("event_id = ?", eventID).First(&event).Error err := db.WithContext(ctx).Where("event_id = ?", eventID).First(&event).Error
if err == nil { if err == nil {
if event.Status != constants.OutboxStatusFailed { if event.Status == constants.OutboxStatusPending || event.Status == constants.OutboxStatusDelivering {
stats.Unchanged++ stats.Unchanged++
return return
} }
result := db.WithContext(ctx).Model(&model.OutboxEvent{}).Where("id = ? AND status = ?", event.ID, constants.OutboxStatusFailed).Updates(map[string]any{ // 已投递只代表入队成功,不代表业务处理成功;业务幂等的退款后处理允许重投。
if event.Status == constants.OutboxStatusDelivered && !retryDelivered {
stats.Unchanged++
return
}
result := db.WithContext(ctx).Model(&model.OutboxEvent{}).Where("id = ? AND status = ?", event.ID, event.Status).Updates(map[string]any{
"status": constants.OutboxStatusPending, "retry_count": 0, "next_attempt_at": time.Now().UTC(), "status": constants.OutboxStatusPending, "retry_count": 0, "next_attempt_at": time.Now().UTC(),
"last_error_code": "", "last_error_summary": "", "updated_at": time.Now().UTC(), "last_error_code": "", "last_error_summary": "", "updated_at": time.Now().UTC(),
}) })

View File

@@ -835,8 +835,8 @@ func (s *Service) deductAllCommission(ctx context.Context, refundID uint) {
} }
} }
// deductSingleCommission 扣单条佣金记录对应的代理钱包余额 // deductSingleCommission 扣单条佣金记录:先拒绝该店铺待审核提现释放冻结余额
// 使用乐观锁扣减(允许余额为负),并创建交易流水 // 再扣减佣金钱包(允许余额为负),最后失效佣金记录并创建流水
func (s *Service) deductSingleCommission(ctx context.Context, refund *model.RefundRequest, commission *model.CommissionRecord) error { func (s *Service) deductSingleCommission(ctx context.Context, refund *model.RefundRequest, commission *model.CommissionRecord) error {
return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var current model.CommissionRecord var current model.CommissionRecord
@@ -869,6 +869,9 @@ func (s *Service) deductSingleCommission(ctx context.Context, refund *model.Refu
First(&wallet).Error; err != nil { First(&wallet).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "锁定退款佣金钱包失败") return errors.Wrap(errors.CodeDatabaseError, err, "锁定退款佣金钱包失败")
} }
if err := s.rejectPendingWithdrawals(ctx, tx, &wallet, current.ShopID, refund); err != nil {
return err
}
result := tx.WithContext(ctx).Model(&model.AgentWallet{}). result := tx.WithContext(ctx).Model(&model.AgentWallet{}).
Where("id = ? AND version = ?", wallet.ID, wallet.Version). Where("id = ? AND version = ?", wallet.ID, wallet.Version).
Updates(map[string]any{"balance": gorm.Expr("balance - ?", current.Amount), "version": gorm.Expr("version + 1")}) Updates(map[string]any{"balance": gorm.Expr("balance - ?", current.Amount), "version": gorm.Expr("version + 1")})
@@ -903,6 +906,97 @@ func (s *Service) deductSingleCommission(ctx context.Context, refund *model.Refu
}) })
} }
// rejectPendingWithdrawals 回扣佣金前拒绝该店铺所有待审核提现。
// 提现冻结的是佣金余额,退款回扣优先级更高;先解冻并拒绝,避免已回扣佣金仍被提现。
func (s *Service) rejectPendingWithdrawals(ctx context.Context, tx *gorm.DB, wallet *model.AgentWallet, shopID uint, refund *model.RefundRequest) error {
var withdrawals []model.CommissionWithdrawalRequest
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).
Where("shop_id = ? AND status = ?", shopID, constants.WithdrawalStatusPending).
Find(&withdrawals).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询待审核佣金提现失败")
}
if len(withdrawals) == 0 {
return nil
}
var shop model.Shop
if err := tx.WithContext(ctx).First(&shop, shopID).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询提现店铺失败")
}
for i := range withdrawals {
w := &withdrawals[i]
before := withdrawalRejectState(w)
if err := s.agentWalletStore.UnfreezeBalanceWithTx(ctx, tx, wallet.ID, w.Amount); err != nil {
return errors.Wrap(errors.CodeInternalError, err, "解冻提现冻结余额失败")
}
refType := constants.ReferenceTypeWithdrawal
remark := "退款佣金回扣,自动拒绝提现"
transaction := &model.AgentWalletTransaction{
AgentWalletID: wallet.ID, ShopID: shopID, UserID: refund.Creator,
TransactionType: constants.AgentTransactionTypeRefund, Amount: w.Amount,
BalanceBefore: wallet.Balance, BalanceAfter: wallet.Balance,
Status: constants.TransactionStatusSuccess, ReferenceType: &refType, ReferenceID: &w.ID,
Remark: &remark, Creator: refund.Creator, ShopIDTag: shopID,
}
if err := tx.WithContext(ctx).Create(transaction).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "创建自动拒绝提现流水失败")
}
now := time.Now()
result := tx.WithContext(ctx).Model(&model.CommissionWithdrawalRequest{}).
Where("id = ? AND status = ?", w.ID, constants.WithdrawalStatusPending).
Updates(map[string]any{
"status": constants.WithdrawalStatusRejected,
"processed_at": now,
"reject_reason": remark,
"updated_at": now,
})
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "拒绝待审核提现失败")
}
if result.RowsAffected != 1 {
return errors.New(errors.CodeConflict, "提现申请状态已变化")
}
w.Status = constants.WithdrawalStatusRejected
w.ProcessedAt = &now
w.RejectReason = remark
if err := s.appendWithdrawalRejectAudit(ctx, tx, w, wallet, transaction, &shop, before); err != nil {
return err
}
}
return nil
}
func withdrawalRejectState(w *model.CommissionWithdrawalRequest) map[string]any {
return map[string]any{
"status": w.Status, "amount": w.Amount, "fee": w.Fee, "actual_amount": w.ActualAmount,
"withdrawal_method": w.WithdrawalMethod, "processed_at": w.ProcessedAt, "reject_reason": w.RejectReason,
}
}
func (s *Service) appendWithdrawalRejectAudit(ctx context.Context, tx *gorm.DB, withdrawal *model.CommissionWithdrawalRequest, wallet *model.AgentWallet, transaction *model.AgentWalletTransaction, shop *model.Shop, before map[string]any) error {
if s.auditWriter == nil {
return errors.New(errors.CodeInvalidStatus, "退款统一审计接缝未配置")
}
primary := audit.CommissionWithdrawalResource(withdrawal, constants.AuditResourceRelationPrimary, constants.AuditResourceRoleWithdrawalTarget,
before, withdrawalRejectState(withdrawal))
primary.SubjectVisibility = constants.AuditSubjectResult
primary.SubjectSummary = "退款佣金回扣自动拒绝提现"
walletResource := audit.AgentWalletResource(wallet, constants.AuditResourceRelationAffected, constants.AuditResourceRoleWithdrawalWallet,
map[string]any{"balance": wallet.Balance, "frozen_balance": wallet.FrozenBalance},
map[string]any{"balance": wallet.Balance, "frozen_balance": wallet.FrozenBalance - withdrawal.Amount})
walletResource.SubjectVisibility = constants.AuditSubjectInternalOnly
transactionResource := audit.AgentWalletTransactionResource(transaction, constants.AuditResourceRelationAffected, constants.AuditResourceRoleWithdrawalTransaction)
transactionResource.SubjectVisibility = constants.AuditSubjectInternalOnly
shopResource := audit.ShopResource(shop, constants.AuditResourceRelationReference, constants.AuditResourceRoleWithdrawalShop)
shopResource.SubjectVisibility = constants.AuditSubjectInternalOnly
return s.auditWriter.Append(ctx, tx, audit.AppendInput{
EventID: "commission-withdrawal:" + strconv.FormatUint(uint64(withdrawal.ID), 10) + ":refund-rejected",
ActionCode: constants.AuditActionCommissionWithdrawalRejected, Summary: "退款佣金回扣自动拒绝提现",
ScopeType: constants.AuditScopePlatform, Result: constants.AuditResultSuccess,
CorrelationID: withdrawal.WithdrawalNo, Metadata: map[string]any{"amount": withdrawal.Amount, "status": withdrawal.Status},
Resources: []audit.ResourceInput{primary, walletResource, transactionResource, shopResource},
})
}
// handleRefundAssetProcessing 幂等处理退款后的资产状态。 // handleRefundAssetProcessing 幂等处理退款后的资产状态。
// 包括退款套餐精准失效、尝试接续待生效主套餐和必要时停机;全部完成后才设置完成标记。 // 包括退款套餐精准失效、尝试接续待生效主套餐和必要时停机;全部完成后才设置完成标记。
func (s *Service) handleRefundAssetProcessing(ctx context.Context, refundID uint) { func (s *Service) handleRefundAssetProcessing(ctx context.Context, refundID uint) {

View File

@@ -0,0 +1,10 @@
-- 回滚到旧的交易类型白名单(不含 commission_deduct 与 adjustment
-- 注意:若已存在 commission_deduct 或 adjustment 流水,回滚会因违反检查约束而失败,属预期保护。
ALTER TABLE tb_agent_wallet_transaction
DROP CONSTRAINT chk_agent_tx_type;
ALTER TABLE tb_agent_wallet_transaction
ADD CONSTRAINT chk_agent_tx_type
CHECK (transaction_type IN ('recharge', 'deduct', 'refund', 'commission', 'withdrawal'));
COMMENT ON COLUMN tb_agent_wallet_transaction.transaction_type IS '交易类型recharge-充值 adjustment-人工调整 deduct-扣款 refund-退款 commission-分佣 withdrawal-提现';

View File

@@ -0,0 +1,20 @@
-- 修复代理钱包交易类型检查约束缺少 commission_deduct 与 adjustment 的问题。
-- 代码常量与列注释早已包含这两种类型,但 chk_agent_tx_type 从未同步,
-- 导致写入退款佣金回扣commission_deduct和人工余额调整adjustment流水时
-- 违反检查约束,退款佣金回扣因此持续失败、无法收回佣金。
ALTER TABLE tb_agent_wallet_transaction
DROP CONSTRAINT chk_agent_tx_type;
ALTER TABLE tb_agent_wallet_transaction
ADD CONSTRAINT chk_agent_tx_type
CHECK (transaction_type IN (
'recharge',
'deduct',
'refund',
'commission',
'withdrawal',
'commission_deduct',
'adjustment'
));
COMMENT ON COLUMN tb_agent_wallet_transaction.transaction_type IS '交易类型recharge-充值 adjustment-人工调整 deduct-扣款 refund-退款 commission-分佣 commission_deduct-退款佣金回扣 withdrawal-提现';

View File

@@ -0,0 +1,8 @@
-- 回滚到旧的交易类型白名单(不含 exchange
-- 注意:若已存在 exchange 流水,回滚会因违反检查约束而失败,属预期保护。
ALTER TABLE tb_asset_wallet_transaction
DROP CONSTRAINT chk_card_tx_type;
ALTER TABLE tb_asset_wallet_transaction
ADD CONSTRAINT chk_card_tx_type
CHECK (transaction_type IN ('recharge', 'deduct', 'refund'));

View File

@@ -0,0 +1,9 @@
-- 修复资产钱包交易类型检查约束缺少 exchange 的问题。
-- 代码常量 AssetTransactionTypeExchange 已用于换货余额迁移,
-- 但 chk_card_tx_type 从未同步,导致写入 exchange 流水时违反检查约束。
ALTER TABLE tb_asset_wallet_transaction
DROP CONSTRAINT chk_card_tx_type;
ALTER TABLE tb_asset_wallet_transaction
ADD CONSTRAINT chk_card_tx_type
CHECK (transaction_type IN ('recharge', 'deduct', 'refund', 'exchange'));

View File

@@ -0,0 +1,13 @@
-- 回滚恢复佣金commission钱包不允许负余额的旧约束。
-- 注意:若已存在负余额佣金数据,回滚会因违反检查约束而失败,属预期保护。
ALTER TABLE tb_agent_wallet
DROP CONSTRAINT chk_agent_wallet_available_balance;
ALTER TABLE tb_agent_wallet
ADD CONSTRAINT chk_agent_wallet_available_balance
CHECK (
(wallet_type = 'main' AND (balance::numeric - frozen_balance::numeric +
CASE WHEN credit_enabled THEN credit_limit::numeric ELSE 0::numeric END) >= 0::numeric)
OR
(wallet_type = 'commission' AND balance >= 0 AND frozen_balance <= balance)
);

View File

@@ -0,0 +1,14 @@
-- 允许佣金commission钱包出现负余额退款佣金回扣优先于提现
-- 回扣后店铺佣金可能倒欠平台,负余额由回扣流程显式允许。
-- 提现冻结仍受 frozen_balance <= GREATEST(balance, 0) 约束:负余额时不允许冻结新提现。
ALTER TABLE tb_agent_wallet
DROP CONSTRAINT chk_agent_wallet_available_balance;
ALTER TABLE tb_agent_wallet
ADD CONSTRAINT chk_agent_wallet_available_balance
CHECK (
(wallet_type = 'main' AND (balance::numeric - frozen_balance::numeric +
CASE WHEN credit_enabled THEN credit_limit::numeric ELSE 0::numeric END) >= 0::numeric)
OR
(wallet_type = 'commission' AND frozen_balance <= GREATEST(balance, 0))
);

View File

@@ -38,6 +38,25 @@
- **WHEN** 操作者为已有系列授权添加套餐 - **WHEN** 操作者为已有系列授权添加套餐
- **THEN** 系统返回同一店铺和系列的候选套餐及其授权状态,且不改变现有套餐管理提交接口的调价和删除语义 - **THEN** 系统返回同一店铺和系列的候选套餐及其授权状态,且不改变现有套餐管理提交接口的调价和删除语义
### Requirement: 套餐使用记录价格快照
系统 SHALL 在创建套餐使用记录时分别快照套餐成本价与零售价:`paid_amount`成本价SHALL 取订单 `seller_cost_price`(销售成本价,即卖家店铺向平台结算的成本),`retail_amount`零售价SHALL 取订单 `total_amount`(零售总价)。成本价与实付金额在个人客户场景下不相等时,`paid_amount` MUST 使用成本价而非实付金额。
#### Scenario: 个人客户购买时快照成本价与零售价
- **WHEN** 个人客户为资产购买套餐,店铺成本价 10900 分,零售价 15900 分,客户实付 15900 分
- **THEN** 创建的套餐使用记录 `paid_amount = 10900``retail_amount = 15900`
#### Scenario: 代理钱包自购时快照成本价与零售价
- **WHEN** 代理以钱包支付为自有资产购买套餐,成本价 7000 分,零售价 9900 分
- **THEN** 创建的套餐使用记录 `paid_amount = 7000``retail_amount = 9900`
#### Scenario: 赠送套餐时快照零成本价与零售价
- **WHEN** 平台赠送套餐,零售价 9900 分,成本价 0 分
- **THEN** 创建的套餐使用记录 `paid_amount = 0``retail_amount = 9900`
## 可达操作索引 ## 可达操作索引
本节只用于入口导航,不是行为 Requirement业务义务以上述 Requirements 为准。 本节只用于入口导航,不是行为 Requirement业务义务以上述 Requirements 为准。