All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m35s
188 lines
8.0 KiB
Go
188 lines
8.0 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{}
|
|
refundStats := 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)
|
|
}
|
|
}
|
|
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)
|
|
}
|
|
if !refund.AssetReset {
|
|
recoverOne(ctx, db, repository, EventRefundAssetProcess, refund.ID, refund.OrderID, &refundStats, 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))
|
|
}
|
|
|
|
func recoverOne(ctx context.Context, db *gorm.DB, repository *outbox.Repository, eventType string, aggregateID, orderID uint, stats *RecoveryStats, logger *zap.Logger) {
|
|
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.OutboxStatusFailed {
|
|
stats.Unchanged++
|
|
return
|
|
}
|
|
result := db.WithContext(ctx).Model(&model.OutboxEvent{}).Where("id = ? AND status = ?", event.ID, constants.OutboxStatusFailed).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++
|
|
}
|