Files
junhong_cmp_fiber/internal/infrastructure/commissiondelivery/event.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

203 lines
8.8 KiB
Go

// Package commissiondelivery 提供订单佣金与退款后处理的可靠 Outbox 事件。
package commissiondelivery
import (
"context"
"strconv"
"time"
"github.com/bytedance/sonic"
"github.com/hibiken/asynq"
"go.uber.org/zap"
"gorm.io/gorm"
"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/outboxid"
)
const (
EventCommissionCalculate = "order.commission.calculate.requested"
EventRefundCommissionDeduct = "refund.commission.deduct.requested"
EventRefundAssetProcess = "refund.asset.process.requested"
PayloadVersionV1 = 1
)
type Payload struct {
OrderID uint `json:"order_id"`
RefundID uint `json:"refund_id,omitempty"`
}
func AppendCommissionCalculate(ctx context.Context, tx *gorm.DB, repository *outbox.Repository, orderID uint) error {
return appendEvent(ctx, tx, repository, EventCommissionCalculate, "order", orderID, Payload{OrderID: orderID})
}
func AppendRefundCommissionDeduct(ctx context.Context, tx *gorm.DB, repository *outbox.Repository, refundID, orderID uint) error {
return appendEvent(ctx, tx, repository, EventRefundCommissionDeduct, "refund", refundID, Payload{OrderID: orderID, RefundID: refundID})
}
func AppendRefundAssetProcess(ctx context.Context, tx *gorm.DB, repository *outbox.Repository, refundID, orderID uint) error {
return appendEvent(ctx, tx, repository, EventRefundAssetProcess, "refund", refundID, Payload{OrderID: orderID, RefundID: refundID})
}
func appendEvent(ctx context.Context, tx *gorm.DB, repository *outbox.Repository, eventType, aggregate string, id uint, payload Payload) error {
if repository == nil {
return gorm.ErrInvalidDB
}
value := strconv.FormatUint(uint64(id), 10)
_, err := repository.AppendIdempotent(ctx, tx, outbox.Envelope{
EventID: outboxid.Stable(eventType+":", value), EventType: eventType, PayloadVersion: PayloadVersionV1,
AggregateType: aggregate, AggregateID: value, ResourceType: aggregate, ResourceID: value,
BusinessKey: eventType + ":" + value, Payload: payload,
})
return err
}
type CommissionConsumer struct {
client outbox.TaskEnqueuer
logger *zap.Logger
}
func NewCommissionConsumer(client outbox.TaskEnqueuer, logger *zap.Logger) *CommissionConsumer {
return &CommissionConsumer{client: client, logger: logger}
}
func (c *CommissionConsumer) Consume(ctx context.Context, envelope outbox.DeliveryEnvelope) error {
var payload Payload
if err := sonic.Unmarshal(envelope.Payload, &payload); err != nil {
return outbox.Permanent(err)
}
if envelope.EventType != EventCommissionCalculate || envelope.PayloadVersion != PayloadVersionV1 || payload.OrderID == 0 {
return outbox.Permanent(gorm.ErrInvalidData)
}
if err := c.client.EnqueueTask(ctx, constants.TaskTypeCommission, map[string]any{"order_id": payload.OrderID, "request_id": envelope.RequestID, "correlation_id": envelope.CorrelationID, "parent_event_id": envelope.EventID}, asynq.Queue(constants.QueueForTaskType(constants.TaskTypeCommission))); err != nil {
return err
}
c.logger.Info("佣金计算 Outbox 已投递", zap.Uint("order_id", payload.OrderID), zap.String("event_id", envelope.EventID), zap.String("correlation_id", envelope.CorrelationID))
return nil
}
type RefundConsumer struct {
commission func(context.Context, uint) error
asset func(context.Context, uint) error
}
func NewRefundConsumer(commission func(context.Context, uint) error, asset func(context.Context, uint) error) *RefundConsumer {
return &RefundConsumer{commission: commission, asset: asset}
}
func (c *RefundConsumer) Consume(ctx context.Context, envelope outbox.DeliveryEnvelope) error {
var payload Payload
if err := sonic.Unmarshal(envelope.Payload, &payload); err != nil {
return outbox.Permanent(err)
}
if payload.RefundID == 0 || payload.OrderID == 0 || envelope.PayloadVersion != PayloadVersionV1 {
return outbox.Permanent(gorm.ErrInvalidData)
}
ctx = auditcontext.With(ctx, auditcontext.Context{CorrelationID: envelope.CorrelationID, ParentEventID: envelope.EventID})
switch envelope.EventType {
case EventRefundCommissionDeduct:
return c.commission(ctx, payload.RefundID)
case EventRefundAssetProcess:
return c.asset(ctx, payload.RefundID)
default:
return outbox.Permanent(gorm.ErrInvalidData)
}
}
// RecoveryStats 是一次有界补偿扫描的可观察结果。
type RecoveryStats struct{ Resent, Unchanged, Failed int }
// Recover 扫描遗留订单与退款,并用稳定事件 ID 恢复投递事实。
func Recover(ctx context.Context, db *gorm.DB, repository *outbox.Repository, limit int, logger *zap.Logger) {
if limit <= 0 {
limit = 100
}
orderStats := RecoveryStats{}
var orders []model.Order
if err := db.WithContext(ctx).Where("payment_status = ? AND commission_status = ?", model.PaymentStatusPaid, model.CommissionStatusPending).Order("id ASC").Limit(limit).Find(&orders).Error; err != nil {
logger.Warn("扫描待计算订单失败", zap.Error(err))
} else {
for _, order := range orders {
recoverOne(ctx, db, repository, EventCommissionCalculate, order.ID, order.ID, &orderStats, logger, false)
}
}
refundStats := RecoverRefundPostProcessing(ctx, db, repository, limit, logger)
logger.Info("佣金与退款补偿扫描完成",
zap.Int("订单已补发", orderStats.Resent), zap.Int("订单无需补发", orderStats.Unchanged), zap.Int("订单失败", orderStats.Failed),
zap.Int("退款已补发", refundStats.Resent), zap.Int("退款无需补发", refundStats.Unchanged), zap.Int("退款失败", refundStats.Failed))
}
// RecoverRefundPostProcessing 只扫描已通过但后处理未闭合的退款单并按稳定业务键恢复缺失或终态失败事件。
// 返还的统计用于启动日志与周期任务日志;本函数只重投事件,不直接改资金。
func RecoverRefundPostProcessing(ctx context.Context, db *gorm.DB, repository *outbox.Repository, limit int, logger *zap.Logger) RecoveryStats {
if limit <= 0 {
limit = 100
}
refundStats := RecoveryStats{}
var refunds []model.RefundRequest
if err := db.WithContext(ctx).Where("status = ? AND (commission_deducted = ? OR asset_reset = ?)", model.RefundStatusApproved, false, false).Order("id ASC").Limit(limit).Find(&refunds).Error; err != nil {
logger.Warn("扫描退款后处理失败", zap.Error(err))
} else {
for _, refund := range refunds {
if !refund.CommissionDeducted {
recoverOne(ctx, db, repository, EventRefundCommissionDeduct, refund.ID, refund.OrderID, &refundStats, logger, true)
}
if !refund.AssetReset {
recoverOne(ctx, db, repository, EventRefundAssetProcess, refund.ID, refund.OrderID, &refundStats, logger, true)
}
}
}
return refundStats
}
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))
var event model.OutboxEvent
err := db.WithContext(ctx).Where("event_id = ?", eventID).First(&event).Error
if err == nil {
if event.Status == constants.OutboxStatusPending || event.Status == constants.OutboxStatusDelivering {
stats.Unchanged++
return
}
// 已投递只代表入队成功,不代表业务处理成功;业务幂等的退款后处理允许重投。
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(),
"last_error_code": "", "last_error_summary": "", "updated_at": time.Now().UTC(),
})
if result.Error == nil && result.RowsAffected == 1 {
stats.Resent++
return
}
if result.Error == nil {
stats.Unchanged++
return
}
stats.Failed++
logger.Warn("恢复失败 Outbox 事件失败", zap.String("event_id", eventID), zap.Error(result.Error))
return
}
if err != gorm.ErrRecordNotFound {
stats.Failed++
logger.Warn("查询补偿 Outbox 事件失败", zap.String("event_id", eventID), zap.Error(err))
return
}
err = db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
switch eventType {
case EventCommissionCalculate:
return AppendCommissionCalculate(ctx, tx, repository, aggregateID)
case EventRefundCommissionDeduct:
return AppendRefundCommissionDeduct(ctx, tx, repository, aggregateID, orderID)
default:
return AppendRefundAssetProcess(ctx, tx, repository, aggregateID, orderID)
}
})
if err != nil {
stats.Failed++
logger.Warn("创建补偿 Outbox 事件失败", zap.String("event_id", eventID), zap.Error(err))
return
}
stats.Resent++
}