修复退款幂等与补偿重试

This commit is contained in:
yml2213
2026-06-14 07:55:02 +08:00
parent 607d2415a1
commit fec8e37181
11 changed files with 772 additions and 143 deletions
+265
View File
@@ -0,0 +1,265 @@
package refundretry
import (
"context"
"crypto/rand"
"encoding/hex"
"fmt"
"time"
"hfb_sys/backend/internal/model"
"hfb_sys/backend/internal/modules/payment"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
"gorm.io/gorm"
)
const (
refundRetryLockKey = "hfb:job:refundretry:lock"
maxRetryCount = 10
baseRetryBackoff = 5 * time.Minute
maxRetryBackoff = time.Hour
manualWarnThreshold = 24 * time.Hour
)
type Job struct {
db *gorm.DB
redis *redis.Client
logger *zap.Logger
payments *payment.Repository
interval time.Duration
instanceID string
}
func New(db *gorm.DB, redisClient *redis.Client, logger *zap.Logger, payments *payment.Repository) *Job {
return &Job{
db: db,
redis: redisClient,
logger: logger,
payments: payments,
interval: 2 * time.Minute,
instanceID: newInstanceID(),
}
}
// acquireLock 通过 Redis 分布式锁确保同一时刻只有一个实例执行退款补偿。
// 未配置 Redis 时直接执行;退款同步本身按退款单精确处理,可容忍短时间重复扫描。
func (j *Job) acquireLock(ctx context.Context) (func(), bool) {
if j.redis == nil {
return func() {}, true
}
ok, err := j.redis.SetNX(ctx, refundRetryLockKey, j.instanceID, j.lockTTL()).Result()
if err != nil {
j.logger.Warn("refund retry job acquire lock failed, run without lock", zap.Error(err))
return func() {}, true
}
if !ok {
return nil, false
}
return j.releaseLock, true
}
// releaseLock 仅在锁仍归本实例时释放,避免误删其他实例的锁。
func (j *Job) releaseLock() {
if j.redis == nil {
return
}
relCtx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
script := redis.NewScript(`if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end`)
if err := script.Run(relCtx, j.redis, []string{refundRetryLockKey}, j.instanceID).Err(); err != nil {
j.logger.Warn("refund retry job release lock failed", zap.Error(err))
}
}
func (j *Job) lockTTL() time.Duration {
return 2 * j.interval
}
func newInstanceID() string {
b := make([]byte, 8)
if _, err := rand.Read(b); err != nil {
return fmt.Sprintf("inst-%d", time.Now().UnixNano())
}
return hex.EncodeToString(b)
}
func (j *Job) Start(ctx context.Context) {
if j == nil || j.db == nil || j.payments == nil {
return
}
go j.loop(ctx)
}
func (j *Job) loop(ctx context.Context) {
j.run(ctx)
ticker := time.NewTicker(j.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
j.logger.Info("refund retry job stopped")
return
case <-ticker.C:
j.run(ctx)
}
}
}
func (j *Job) run(ctx context.Context) {
release, ok := j.acquireLock(ctx)
if !ok {
return
}
defer release()
now := time.Now()
processed, err := j.syncRefundPayments(ctx, now)
if err != nil {
j.logger.Warn("refund retry job sync payments failed", zap.Error(err))
}
missing, err := j.warnMissingRefundOrders(ctx, now)
if err != nil {
j.logger.Warn("refund retry job scan missing refund orders failed", zap.Error(err))
}
if processed > 0 || missing > 0 {
j.logger.Info("refund retry job finished", zap.Int("processed", processed), zap.Int("missing_refund_orders", missing))
}
}
func (j *Job) syncRefundPayments(ctx context.Context, now time.Time) (int, error) {
var rows []model.PaymentOrder
err := j.db.WithContext(ctx).
Where("biz_type IN ? AND status IN ? AND updated_at <= ? AND retry_count < ? AND (next_retry_at IS NULL OR next_retry_at <= ?)",
payment.RefundBizTypes(), []string{"refunding", "failed"}, now.Add(-5*time.Minute), maxRetryCount, now).
Order("id ASC").
Limit(100).
Find(&rows).Error
if err != nil {
return 0, err
}
processed := 0
for _, row := range rows {
if _, err := j.payments.SyncRefundStatusByPaymentID(ctx, row.ID); err != nil {
if markErr := j.markRetryFailed(ctx, row, now); markErr != nil {
j.logger.Warn("refund retry mark failure failed",
zap.Uint64("payment_id", row.ID),
zap.Uint64("order_id", row.OrderID),
zap.Error(markErr),
)
}
j.logger.Warn("refund retry sync failed",
zap.Uint64("payment_id", row.ID),
zap.Uint64("order_id", row.OrderID),
zap.String("biz_type", row.BizType),
zap.String("status", row.Status),
zap.Int("retry_count", row.RetryCount+1),
zap.Error(err),
)
continue
}
if err := j.resetRetry(ctx, row.ID); err != nil {
j.logger.Warn("refund retry reset counter failed",
zap.Uint64("payment_id", row.ID),
zap.Uint64("order_id", row.OrderID),
zap.Error(err),
)
}
processed++
}
j.warnMaxRetryRefunds(ctx, now)
return processed, nil
}
func (j *Job) markRetryFailed(ctx context.Context, row model.PaymentOrder, now time.Time) error {
nextCount := row.RetryCount + 1
backoff := retryBackoff(nextCount)
nextRetryAt := now.Add(backoff)
updates := map[string]any{
"retry_count": nextCount,
"last_retry_at": now,
"next_retry_at": nextRetryAt,
}
if nextCount >= maxRetryCount {
updates["next_retry_at"] = nil
}
return j.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", row.ID).Updates(updates).Error
}
func (j *Job) resetRetry(ctx context.Context, paymentID uint64) error {
return j.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", paymentID).Updates(map[string]any{
"retry_count": 0,
"last_retry_at": nil,
"next_retry_at": nil,
}).Error
}
func retryBackoff(retryCount int) time.Duration {
if retryCount <= 1 {
return baseRetryBackoff
}
backoff := baseRetryBackoff
for i := 1; i < retryCount; i++ {
backoff *= 2
if backoff >= maxRetryBackoff {
return maxRetryBackoff
}
}
return backoff
}
func (j *Job) warnMaxRetryRefunds(ctx context.Context, now time.Time) {
var rows []model.PaymentOrder
err := j.db.WithContext(ctx).
Where("biz_type IN ? AND status IN ? AND retry_count >= ? AND (last_retry_at IS NULL OR last_retry_at <= ?)",
payment.RefundBizTypes(), []string{"refunding", "failed"}, maxRetryCount, now.Add(-manualWarnThreshold)).
Order("id ASC").
Limit(50).
Find(&rows).Error
if err != nil {
j.logger.Warn("refund retry max count scan failed", zap.Error(err))
return
}
for _, row := range rows {
j.logger.Warn("refund retry reached max count, need manual check",
zap.Uint64("payment_id", row.ID),
zap.Uint64("order_id", row.OrderID),
zap.String("biz_type", row.BizType),
zap.String("status", row.Status),
zap.Int("retry_count", row.RetryCount),
)
}
}
func (j *Job) warnMissingRefundOrders(ctx context.Context, now time.Time) (int, error) {
var rows []model.RentalOrder
err := j.db.WithContext(ctx).
Where("refund_status IN ? AND refund_amount_cent > 0 AND updated_at <= ?", []string{"pending", "refunding"}, now.Add(-10*time.Minute)).
Order("id ASC").
Limit(100).
Find(&rows).Error
if err != nil {
return 0, err
}
missing := 0
for _, row := range rows {
var count int64
if err := j.db.WithContext(ctx).Model(&model.PaymentOrder{}).
Where("order_id = ? AND biz_type IN ?", row.ID, payment.RefundBizTypes()).
Count(&count).Error; err != nil {
return missing, err
}
if count > 0 {
continue
}
missing++
j.logger.Warn("refund order missing payment record, need manual check",
zap.Uint64("order_id", row.ID),
zap.String("order_no", row.OrderNo),
zap.String("refund_status", row.RefundStatus),
zap.Int64("refund_amount_cent", row.RefundAmountCent),
)
}
return missing, nil
}
@@ -0,0 +1,117 @@
package refundretry
import (
"testing"
"time"
"hfb_sys/backend/internal/model"
"go.uber.org/zap"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"gorm.io/gorm/logger"
)
func setupRefundRetryTestDB(t *testing.T) *gorm.DB {
t.Helper()
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
Logger: logger.Default.LogMode(logger.Silent),
})
if err != nil {
t.Fatalf("创建测试数据库失败: %v", err)
}
if err := db.AutoMigrate(&model.PaymentOrder{}, &model.RentalOrder{}); err != nil {
t.Fatalf("数据库迁移失败: %v", err)
}
return db
}
func TestRetryBackoffCapsAtOneHour(t *testing.T) {
cases := []struct {
retryCount int
want time.Duration
}{
{retryCount: 1, want: 5 * time.Minute},
{retryCount: 2, want: 10 * time.Minute},
{retryCount: 3, want: 20 * time.Minute},
{retryCount: 4, want: 40 * time.Minute},
{retryCount: 5, want: time.Hour},
{retryCount: 10, want: time.Hour},
}
for _, tc := range cases {
if got := retryBackoff(tc.retryCount); got != tc.want {
t.Fatalf("retryBackoff(%d) = %s, want %s", tc.retryCount, got, tc.want)
}
}
}
func TestMarkRetryFailedStopsAtMaxCount(t *testing.T) {
db := setupRefundRetryTestDB(t)
job := New(db, nil, zap.NewNop(), nil)
now := time.Date(2026, 6, 14, 12, 0, 0, 0, time.UTC)
row := model.PaymentOrder{
PaymentNo: "PAY202606140101",
OrderID: 1,
OrderNo: "ORD202606140101",
UserID: 2,
Provider: "mock",
ThirdOrderID: "REF202606140101",
BizType: "admin_refund",
Status: "failed",
RetryCount: maxRetryCount - 1,
}
if err := db.Create(&row).Error; err != nil {
t.Fatalf("创建退款单失败: %v", err)
}
if err := job.markRetryFailed(t.Context(), row, now); err != nil {
t.Fatalf("markRetryFailed() error = %v", err)
}
var latest model.PaymentOrder
if err := db.First(&latest, row.ID).Error; err != nil {
t.Fatalf("查询退款单失败: %v", err)
}
if latest.RetryCount != maxRetryCount {
t.Fatalf("RetryCount = %d, want %d", latest.RetryCount, maxRetryCount)
}
if latest.LastRetryAt == nil || !latest.LastRetryAt.Equal(now) {
t.Fatalf("LastRetryAt = %v, want %v", latest.LastRetryAt, now)
}
if latest.NextRetryAt != nil {
t.Fatalf("NextRetryAt = %v, want nil after max retry", latest.NextRetryAt)
}
}
func TestResetRetryClearsRetryFields(t *testing.T) {
db := setupRefundRetryTestDB(t)
job := New(db, nil, zap.NewNop(), nil)
now := time.Date(2026, 6, 14, 12, 0, 0, 0, time.UTC)
next := now.Add(time.Hour)
row := model.PaymentOrder{
PaymentNo: "PAY202606140102",
OrderID: 1,
OrderNo: "ORD202606140102",
UserID: 2,
Provider: "mock",
ThirdOrderID: "REF202606140102",
BizType: "admin_refund",
Status: "refunded",
RetryCount: 3,
LastRetryAt: &now,
NextRetryAt: &next,
}
if err := db.Create(&row).Error; err != nil {
t.Fatalf("创建退款单失败: %v", err)
}
if err := job.resetRetry(t.Context(), row.ID); err != nil {
t.Fatalf("resetRetry() error = %v", err)
}
var latest model.PaymentOrder
if err := db.First(&latest, row.ID).Error; err != nil {
t.Fatalf("查询退款单失败: %v", err)
}
if latest.RetryCount != 0 || latest.LastRetryAt != nil || latest.NextRetryAt != nil {
t.Fatalf("retry fields = count:%d last:%v next:%v, want zero/nil", latest.RetryCount, latest.LastRetryAt, latest.NextRetryAt)
}
}