Files
junhong_cmp_fiber/internal/service/refund/approval_decision.go
break 1aa4eacee2
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 9m26s
feat(退款分佣): 佣金回溯明细替换全额失效并补齐读侧与导出
用 PRD 2.14 语义整体替换退款佣金「整单全额失效」实现:原佣金保持已发放不变,
回溯事实落在新表 tb_commission_clawback_record 的负数、不可提现明细上。

- 新增成对迁移 000220 建 tb_commission_clawback_record,唯一约束
  (refund_id, original_commission_id) 为权威幂等键,附店铺+时间/原佣金/订单索引。
- 回溯用例(internal/service/refund/clawback.go):准入仅由退款申请状态、审批异常
  标记与退款方式决定;金额按分整数计算,分母取冻结实收(缺失回落审批尝试)、
  分子原路取渠道成功金额,乘法用 math/big 中间量,舍入差自末条起向前补差;
  终态判据要求订单佣金已离开待计算且不存在 status IN (1,2,99) 的记录。
- 三层幂等:唯一约束兜底、佣金行行锁 + 钱包乐观锁、commission_deducted 仅作投影
  并带 WHERE commission_deducted = false 条件置位;闭合三结果为已回溯、无需回溯、
  审批异常转人工。
- 事务内顺序固定:锁提现申请行 → 锁尝试行 → 解冻冻结 → 置驳回 → 插回溯明细 →
  扣 balance(允许为负)→ 写负数流水 → 审计;删除旧全额失效写入与其两个审计调用点,
  refund.invalidate_commission 仅保留常量与注册供历史审计读取。
- 读侧:佣金明细列表 status 筛选透传,两表 UNION ALL 合并分页并以 source ASC 作
  末位次序键;新增佣金明细详情接口并同步路由与 OpenAPI 装配。
- 导出:新增 commission_record 场景(白名单、exporter 注册、DTO oneof、DataSource
  与列定义),粒度为佣金记录,原佣金与回溯各一行,金额保持分且可为负。
- 新增退款佣金回溯周期补偿任务(@every 1m / MaxRetry(3) / Timeout(10m) /
  Unique(10m),独立队列),保留启动时补偿扫描,判据与既有实现一致。

Refs: AUG26-012
2026-09-14 13:40:34 +08:00

284 lines
13 KiB
Go
Raw Permalink 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 refund
import (
"context"
"strconv"
"strings"
"gorm.io/gorm"
"gorm.io/gorm/clause"
approvalapp "github.com/break/junhong_cmp_fiber/internal/application/approval"
employeecollectionapp "github.com/break/junhong_cmp_fiber/internal/application/employeecollection"
refundapproval "github.com/break/junhong_cmp_fiber/internal/application/refundapproval"
"github.com/break/junhong_cmp_fiber/internal/infrastructure/commissiondelivery"
"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/auditcontext"
"github.com/break/junhong_cmp_fiber/pkg/constants"
"github.com/break/junhong_cmp_fiber/pkg/errors"
)
// Handle 幂等处理退款的渠道无关审批终态。
func (s *Service) Handle(ctx context.Context, event approvalapp.TerminalDecisionEvent) error {
if s == nil || s.db == nil || event.BusinessType != constants.ApprovalBusinessTypeRefund ||
event.BusinessID == 0 || event.InstanceID == 0 {
return errors.New(errors.CodeInvalidParam, "退款审批终态参数无效")
}
switch event.Decision {
case constants.ApprovalDecisionApproved:
return s.applyApprovedDecision(ctx, event)
case constants.ApprovalDecisionRejected, constants.ApprovalDecisionCancelled, constants.ApprovalDecisionDeleted:
return s.applyClosedDecision(ctx, event)
case constants.ApprovalDecisionRevokedAfterApproved:
return s.applyRevokedAfterApproved(ctx, event)
default:
return errors.New(errors.CodeInvalidParam, "不支持的退款审批终态")
}
}
// applyRevokedAfterApproved 处理企业微信通过后撤销。
//
// 不回滚已失效的套餐权益、不取消已提交的渠道退款、也不恢复订单退款状态:这些事实在通过时
// 已经成立。系统只标记审批异常并禁止后续自动重提,交由超级管理员线下处理。
func (s *Service) applyRevokedAfterApproved(ctx context.Context, event approvalapp.TerminalDecisionEvent) error {
ctx = auditcontext.With(ctx, auditcontext.Context{CorrelationID: event.CorrelationID, ParentEventID: event.EventID})
var refund model.RefundRequest
err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
target, _, err := refundapproval.ResolveRefundInTx(ctx, tx, event.BusinessID, event.InstanceID)
if err != nil {
return err
}
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).First(&refund, target.ID).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "锁定退款申请失败")
}
if refund.AnomalyFlag == 1 {
return nil
}
beforeRefund := refundAuditState(&refund)
result := tx.WithContext(ctx).Model(&model.RefundRequest{}).
Where("id = ? AND anomaly_flag = 0", refund.ID).
Updates(map[string]any{
"anomaly_flag": 1,
"anomaly_reason": "企业微信通过后撤销,需超级管理员线下处理",
"failure_reason": constants.RefundFailureRevokedAfterApproved,
"updated_at": event.OccurredAt,
})
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "标记退款审批异常失败")
}
if result.RowsAffected != 1 {
return nil
}
refund.AnomalyFlag = 1
refund.AnomalyReason = "企业微信通过后撤销,需超级管理员线下处理"
refund.FailureReason = constants.RefundFailureRevokedAfterApproved
return s.appendRefundAudit(ctx, tx, refund.ID, constants.AuditActionRefundAnomalyFlagged, "标记退款审批异常",
"refund:"+strconv.FormatUint(uint64(refund.ID), 10)+":revoked", beforeRefund, nil, "退款审批通过后撤销")
})
if err != nil {
return err
}
return nil
}
func (s *Service) applyApprovedDecision(ctx context.Context, event approvalapp.TerminalDecisionEvent) error {
ctx = auditcontext.With(ctx, auditcontext.Context{CorrelationID: event.CorrelationID, ParentEventID: event.EventID})
var refund model.RefundRequest
var order model.Order
changed := false
err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
// 业务标识解析:审批尝试记录优先,退款申请兜底(兼容尚未接入尝试模式的存量申请)。
target, attempt, err := refundapproval.ResolveRefundInTx(ctx, tx, event.BusinessID, event.InstanceID)
if err != nil {
return err
}
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).First(&refund, target.ID).Error; err != nil {
if err == gorm.ErrRecordNotFound {
return errors.New(errors.CodeForbidden, "无权限操作该资源或资源不存在")
}
return errors.Wrap(errors.CodeDatabaseError, err, "锁定退款申请失败")
}
if refund.Status != model.RefundStatusPending && refund.Status != model.RefundStatusApproved &&
refund.Status != model.RefundStatusChannelProcessing {
return errors.New(errors.CodeInvalidStatus, "退款申请状态不允许审批通过")
}
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", refund.OrderID).First(&order).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "锁定退款关联订单失败")
}
beforeRefund := refundAuditState(&refund)
beforeOrder := map[string]any{"payment_status": order.PaymentStatus}
approvedAmount := refund.RequestedRefundAmount
if err := validateApprovedRefundAmount(approvedAmount, refund.RequestedRefundAmount, &order); err != nil {
return err
}
if err := s.preparePaymentRefundCredentials(ctx, tx, order.ID); err != nil {
return err
}
// 原路退款只登记待执行事实:退款单转入原路处理中,订单在渠道明确成功前保持已支付。
// 客户收款信息与退回原钱包在审批通过时即完成,订单同时置为已退款。
originalRoute := refund.Method == constants.RefundMethodOriginalRoute
if refund.Status == model.RefundStatusPending {
changed = true
targetStatus := model.RefundStatusApproved
if originalRoute {
targetStatus = model.RefundStatusChannelProcessing
}
result := tx.WithContext(ctx).Model(&model.RefundRequest{}).
Where("id = ? AND status = ?", refund.ID, model.RefundStatusPending).
Updates(map[string]any{
"status": targetStatus, "processed_at": event.OccurredAt,
"approved_refund_amount": approvedAmount, "remark": "企业微信审批通过",
"failure_reason": "", "failure_message": "",
"updated_at": event.OccurredAt,
})
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "完成退款审批申请失败")
}
if result.RowsAffected != 1 {
return errors.New(errors.CodeConflict, "退款申请状态已变化")
}
refund.Status = targetStatus
}
if !originalRoute {
switch order.PaymentStatus {
case model.PaymentStatusPaid:
changed = true
result := tx.WithContext(ctx).Model(&model.Order{}).
Where("id = ? AND payment_status = ?", order.ID, model.PaymentStatusPaid).
Updates(map[string]any{"payment_status": model.PaymentStatusRefunded, "updated_at": event.OccurredAt})
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "更新订单退款状态失败")
}
if result.RowsAffected != 1 {
return errors.New(errors.CodeConflict, "订单退款状态已变化")
}
order.PaymentStatus = model.PaymentStatusRefunded
case model.PaymentStatusRefunded:
default:
return errors.New(errors.CodeInvalidStatus, "订单状态不允许完成退款")
}
} else {
if order.PaymentStatus != model.PaymentStatusPaid && order.PaymentStatus != model.PaymentStatusRefunded {
return errors.New(errors.CodeInvalidStatus, "订单状态不允许原路退款")
}
if err := s.prepareChannelRefundInTx(ctx, tx, &refund, attempt); err != nil {
return err
}
}
if err := s.refundWalletPayment(ctx, tx, &refund, &order, approvedAmount, event.SubmitterAccountID); err != nil {
return err
}
// 员工代收款冲销:接入点只在既有退款成功事务内,以 bill_id+refund_id 唯一事实幂等,
// 不依赖下方 changed 门;来源订单未建账时直接跳过,不阻断退款。
if s.refundOffset == nil {
return errors.New(errors.CodeInternalError, "员工代收款退款冲销能力未配置")
}
if err := s.refundOffset.ApplyInTx(ctx, tx, employeecollectionapp.RefundOffsetSource{
RefundID: refund.ID, OrderID: order.ID, RefundAmount: approvedAmount,
}); err != nil {
return err
}
if !originalRoute {
if err := s.appendCompletedNotification(ctx, tx, &refund); err != nil {
return err
}
}
if err := commissiondelivery.AppendRefundCommissionDeduct(ctx, tx, outbox.NewRepository(), refund.ID, refund.OrderID); err != nil {
return err
}
if err := commissiondelivery.AppendRefundAssetProcess(ctx, tx, outbox.NewRepository(), refund.ID, refund.OrderID); err != nil {
return err
}
if !changed {
return nil
}
return s.appendRefundAudit(ctx, tx, refund.ID, constants.AuditActionRefundApproved, "通过退款审批",
"refund:"+strconv.FormatUint(uint64(refund.ID), 10)+":approved", beforeRefund, beforeOrder, "退款已通过")
})
if err != nil {
s.recordRefundFailure(ctx, constants.AuditActionRefundApproved, "通过退款审批失败", &refund, &order, err)
return err
}
return nil
}
// prepareChannelRefundInTx 在企微通过事务内登记原路退款的待执行事实。
//
// 只写本地事实与可靠事件:请求号在提交时已冻结到审批尝试记录,这里把它落到退款单并转入
// 原路处理中真正的渠道调用由可靠事件驱动的事务外消费者执行ENG-TX-001
func (s *Service) prepareChannelRefundInTx(ctx context.Context, tx *gorm.DB, refund *model.RefundRequest, attempt *model.RefundRequestAttempt) error {
if s.channelRefund == nil {
return errors.New(errors.CodeInternalError, "渠道原路退款能力未配置")
}
return s.channelRefund.PrepareInTx(ctx, tx, refund, attempt)
}
func (s *Service) applyClosedDecision(ctx context.Context, event approvalapp.TerminalDecisionEvent) error {
ctx = auditcontext.With(ctx, auditcontext.Context{CorrelationID: event.CorrelationID, ParentEventID: event.EventID})
reason := map[string]string{
constants.ApprovalDecisionRejected: "企业微信审批已拒绝",
constants.ApprovalDecisionCancelled: "企业微信审批已撤销",
constants.ApprovalDecisionDeleted: "企业微信审批已删除",
}[event.Decision]
var refund model.RefundRequest
var order model.Order
err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.WithContext(ctx).Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", event.BusinessID).First(&refund).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "锁定退款申请失败")
}
if refund.ApprovalInstanceID == nil || *refund.ApprovalInstanceID != event.InstanceID {
return errors.New(errors.CodeConflict, "退款申请关联的审批实例不一致")
}
if refund.Status == model.RefundStatusRejected {
return nil
}
if refund.Status != model.RefundStatusPending {
return errors.New(errors.CodeInvalidStatus, "退款申请状态不允许结束审批")
}
if err := tx.WithContext(ctx).First(&order, refund.OrderID).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询退款关联订单失败")
}
beforeRefund := refundAuditState(&refund)
result := tx.WithContext(ctx).Model(&model.RefundRequest{}).
Where("id = ? AND status = ?", refund.ID, model.RefundStatusPending).
Updates(map[string]any{
"status": model.RefundStatusRejected, "processed_at": event.OccurredAt,
"reject_reason": strings.TrimSpace(reason), "updated_at": event.OccurredAt,
})
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "结束退款审批申请失败")
}
if result.RowsAffected != 1 {
return errors.New(errors.CodeConflict, "退款申请状态已变化")
}
return s.appendRefundAudit(ctx, tx, refund.ID, constants.AuditActionRefundRejected, "拒绝退款审批",
"refund:"+strconv.FormatUint(uint64(refund.ID), 10)+":rejected", beforeRefund, nil, "退款已拒绝")
})
if err != nil {
s.recordRefundFailure(ctx, constants.AuditActionRefundRejected, "拒绝退款审批失败", &refund, &order, err)
}
return err
}
// ProcessCommissionDeduction 处理退款佣金回溯后处理。
// 返回非 nil 表示后处理尚未闭合(准入未满足、冻结实收非正或订单佣金未终态),
// 由 Outbox 与周期性补偿任务重试;返回 nil 表示已闭合并已落库全部应有事实。
func (s *Service) ProcessCommissionDeduction(ctx context.Context, refundID uint) error {
return s.clawbackCommission(ctx, refundID)
}
func (s *Service) ProcessAssetPostProcessing(ctx context.Context, refundID uint) error {
s.handleRefundAssetProcessing(ctx, refundID)
var refund model.RefundRequest
if err := s.db.WithContext(ctx).Select("asset_reset").First(&refund, refundID).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "复核退款资产后处理状态失败")
}
if !refund.AssetReset {
return errors.New(errors.CodeServiceUnavailable, "退款资产后处理尚未完成")
}
return nil
}
var _ approvalapp.BusinessDecisionHandler = (*Service)(nil)