diff --git a/backend/cmd/api/main.go b/backend/cmd/api/main.go index d9b1dd5..ef7c0ed 100644 --- a/backend/cmd/api/main.go +++ b/backend/cmd/api/main.go @@ -11,11 +11,15 @@ import ( "hfb_sys/backend/internal/config" "hfb_sys/backend/internal/database" "hfb_sys/backend/internal/jobs/ordertimeout" + "hfb_sys/backend/internal/jobs/refundretry" "hfb_sys/backend/internal/logging" "hfb_sys/backend/internal/modules/adminauth" + "hfb_sys/backend/internal/modules/payment" + "hfb_sys/backend/internal/modules/paymentconfig" "hfb_sys/backend/internal/router" "go.uber.org/zap" + "gorm.io/gorm" ) // @title HFB Sys API @@ -84,6 +88,10 @@ func main() { defer stopJobs() if deps.DB != nil { ordertimeout.New(deps.DB, deps.Redis, logger).Start(jobCtx) + if paymentConfigRepo := newPaymentConfigRepositoryForJobs(cfg, deps.DB, logger); paymentConfigRepo != nil { + paymentRepo := payment.NewRepository(deps.DB, paymentConfigRepo, nil) + refundretry.New(deps.DB, deps.Redis, logger, paymentRepo).Start(jobCtx) + } } server := &http.Server{ Addr: cfg.AppAddr, @@ -110,3 +118,21 @@ func main() { } logger.Info("api server stopped") } + +func newPaymentConfigRepositoryForJobs(cfg config.Config, db *gorm.DB, logger *zap.Logger) *paymentconfig.Repository { + encryptionKey := cfg.PaymentConfigEncryptionKey + var encryptor paymentconfig.Encryptor + if encryptionKey != "" { + if aesEncryptor, err := paymentconfig.NewAESEncryptor(encryptionKey); err == nil { + encryptor = aesEncryptor + } + } + if encryptor == nil { + if cfg.AppEnv == "production" { + logger.Fatal("PAYMENT_CONFIG_ENCRYPTION_KEY not set or invalid") + } + encryptor = &paymentconfig.MockEncryptor{} + logger.Warn("PAYMENT_CONFIG_ENCRYPTION_KEY not set or invalid, using MockEncryptor for jobs") + } + return paymentconfig.NewRepository(db, encryptor) +} diff --git a/backend/internal/jobs/refundretry/job.go b/backend/internal/jobs/refundretry/job.go new file mode 100644 index 0000000..cba7d82 --- /dev/null +++ b/backend/internal/jobs/refundretry/job.go @@ -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 +} diff --git a/backend/internal/jobs/refundretry/job_test.go b/backend/internal/jobs/refundretry/job_test.go new file mode 100644 index 0000000..9b7b118 --- /dev/null +++ b/backend/internal/jobs/refundretry/job_test.go @@ -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) + } +} diff --git a/backend/internal/model/payment.go b/backend/internal/model/payment.go index 465dc1a..3ffcefc 100644 --- a/backend/internal/model/payment.go +++ b/backend/internal/model/payment.go @@ -26,6 +26,9 @@ type PaymentOrder struct { JSPayInfo string `gorm:"column:jspay_info;type:text" json:"jspay_info"` RawRequest datatypes.JSON `json:"raw_request"` RawResponse datatypes.JSON `json:"raw_response"` + RetryCount int `gorm:"not null;default:0" json:"retry_count"` + LastRetryAt *time.Time `json:"last_retry_at"` + NextRetryAt *time.Time `json:"next_retry_at"` PaidAt *time.Time `json:"paid_at"` NotifiedAt *time.Time `json:"notified_at"` CreatedAt time.Time `json:"created_at"` diff --git a/backend/internal/modules/dispute/arbitration.go b/backend/internal/modules/dispute/arbitration.go index c91bc0f..e730b5c 100644 --- a/backend/internal/modules/dispute/arbitration.go +++ b/backend/internal/modules/dispute/arbitration.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "log" "time" "hfb_sys/backend/internal/model" @@ -274,7 +275,9 @@ func (r *Repository) startRefundBestEffort(ctx context.Context, action *refundAc if action == nil || r.refundStarter == nil { return } - _, _ = r.refundStarter.StartRefund(ctx, action.OrderID, action.RefundAmountCent, action.BizType, action.Remark) + if _, err := r.refundStarter.StartRefund(ctx, action.OrderID, action.RefundAmountCent, action.BizType, action.Remark); err != nil { + log.Printf("[dispute] start refund failed order_id=%d biz_type=%s amount_cent=%d err=%v", action.OrderID, action.BizType, action.RefundAmountCent, err) + } } func renterFrozenBalance(tx *gorm.DB, renterID uint64) (int64, error) { diff --git a/backend/internal/modules/payment/notify.go b/backend/internal/modules/payment/notify.go index e5f3326..1be0e91 100644 --- a/backend/internal/modules/payment/notify.go +++ b/backend/internal/modules/payment/notify.go @@ -5,7 +5,6 @@ import ( "gorm.io/gorm" "hfb_sys/backend/internal/model" "log" - "time" ) func (r *Repository) HandleLeshuaNotify(ctx context.Context, params map[string]string, rawPayload string, contentType string) (*NotifyResult, error) { @@ -58,30 +57,14 @@ func (r *Repository) HandleRefundNotify(ctx context.Context, provider string, pa raw := withNotifyDiagnostic(params, rawPayload, contentType, verify, "verified") status := normalizeNotifyRefundStatus(provider, params["status"]) - switch status { - case "refunded": - now := time.Now() - if err := r.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", payment.ID).Updates(map[string]any{ - "status": "refunded", - "paid_at": now, - "notified_at": now, - "raw_response": jsonMap(raw), - }).Error; err != nil { - return nil, err - } - _ = r.updateOrderRefundStatus(ctx, payment.OrderID, payment.AmountCent) - case "failed": - r.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", payment.ID).Updates(map[string]any{ - "status": "failed", - "notified_at": time.Now(), - "raw_response": jsonMap(raw), - }) - _ = r.markOrderRefundFailed(ctx, payment.OrderID, payment.AmountCent) - default: - r.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", payment.ID).Updates(map[string]any{ - "status": "refunding", - "raw_response": jsonMap(raw), - }) + if err := r.applyRefundChannelStatus(ctx, payment, refundChannelStatusUpdate{ + Status: status, + ProviderRefundID: firstNonEmpty(params["provider_refund_id"], params["leshua_refund_id"]), + RefundTime: params["refund_time"], + Raw: raw, + Source: channelSourceNotify, + }); err != nil { + return nil, err } return &NotifyResult{OK: true, Message: "000000"}, nil } diff --git a/backend/internal/modules/payment/refund.go b/backend/internal/modules/payment/refund.go index a425d7f..e8f55b2 100644 --- a/backend/internal/modules/payment/refund.go +++ b/backend/internal/modules/payment/refund.go @@ -3,21 +3,18 @@ package payment import ( "context" "encoding/json" - "fmt" "log" "time" "gorm.io/datatypes" "gorm.io/gorm" + "gorm.io/gorm/clause" "hfb_sys/backend/internal/model" ) func (r *Repository) StartRefund(ctx context.Context, orderID uint64, refundAmountCent int64, bizType string, remark string) (*RefundDTO, error) { - var originalPayment model.PaymentOrder - if err := r.db.WithContext(ctx).Where("order_id = ? AND status = 'paid' AND biz_type = 'order_pay'", orderID).Order("id DESC").First(&originalPayment).Error; err != nil { - if err == gorm.ErrRecordNotFound { - return nil, ErrPaymentNotFound - } + originalPayment, err := r.findOriginalPayment(ctx, orderID) + if err != nil { return nil, err } runtimeConfig, err := r.runtimeConfigForPayment(ctx, &originalPayment) @@ -25,126 +22,99 @@ func (r *Repository) StartRefund(ctx context.Context, orderID uint64, refundAmou return nil, ErrPaymentUnavailable } - var existingRefund model.PaymentOrder - err = r.db.WithContext(ctx).Where("order_id = ? AND biz_type = ? AND status NOT IN ('failed')", orderID, bizType).Order("id DESC").First(&existingRefund).Error - if err == nil { - dto := toRefundDTO(existingRefund) - return &dto, nil - } - if err != gorm.ErrRecordNotFound { - return nil, err - } - - paymentNo, err := newPaymentNo() + refundOrder, existing, err := r.prepareRefundOrder(ctx, originalPayment, *runtimeConfig, refundAmountCent, bizType) if err != nil { return nil, err } - merchantRefundID := "REF" + paymentNo[3:] - - refundOrder := model.PaymentOrder{ - PaymentNo: paymentNo, - OrderID: orderID, - OrderNo: originalPayment.OrderNo, - UserID: originalPayment.UserID, - Provider: runtimeConfig.Provider, - MerchantID: runtimeConfig.MerchantID, - ThirdOrderID: merchantRefundID, - ProviderOrderID: "", - PayWay: originalPayment.PayWay, - JSPayFlag: originalPayment.JSPayFlag, - AmountCent: refundAmountCent, - BizType: bizType, - Status: "refunding", - } - - if runtimeConfig.isMockMode() { - refundOrder.ProviderOrderID = "MOCKREF" + merchantRefundID - refundOrder.Status = "refunded" - now := time.Now() - refundOrder.PaidAt = &now - if remark != "" { - refundOrder.RawResponse = datatypes.JSON([]byte(fmt.Sprintf(`{"mock":"true","remark":"%s"}`, remark))) + if existing { + latest, syncErr := r.syncRefundPayment(ctx, refundOrder, refundOrder.Status != "refunded") + if syncErr != nil { + log.Printf("[payment] sync existing refund failed order_id=%d payment_id=%d biz_type=%s err=%v", orderID, refundOrder.ID, bizType, syncErr) + dto := toRefundDTO(*refundOrder) + return &dto, nil } - if err := r.db.WithContext(ctx).Create(&refundOrder).Error; err != nil { - return nil, err - } - r.recordConfigUsage(ctx, runtimeConfig, &refundOrder) - if err := r.updateOrderRefundStatus(ctx, orderID, refundAmountCent); err != nil { - log.Printf("[payment] mock update order refund status failed order_id=%d err=%v", orderID, err) - } - dto := toRefundDTO(refundOrder) + dto := toRefundDTO(*latest) return &dto, nil } - if err := r.db.WithContext(ctx).Create(&refundOrder).Error; err != nil { - return nil, err + if runtimeConfig.isMockMode() { + raw := map[string]string{"mock": "true"} + if remark != "" { + raw["remark"] = remark + } + if err := r.applyRefundChannelStatus(ctx, refundOrder, refundChannelStatusUpdate{ + Status: "refunded", + ProviderRefundID: "MOCKREF" + refundOrder.ThirdOrderID, + Raw: raw, + Source: channelSourceMock, + }); err != nil { + return nil, err + } + latest, err := r.findPaymentByID(ctx, refundOrder.ID) + if err != nil { + return nil, err + } + r.recordConfigUsage(ctx, runtimeConfig, latest) + dto := toRefundDTO(*latest) + return &dto, nil } + log.Printf("[payment] refund start order_id=%d order_no=%s payment_id=%d biz_type=%s provider=%s amount_cent=%d merchant_refund_id=%s origin_third_order_id=%s origin_provider_order_id=%s", - orderID, originalPayment.OrderNo, refundOrder.ID, bizType, runtimeConfig.Provider, refundAmountCent, merchantRefundID, originalPayment.ThirdOrderID, refundOriginProviderOrderID(originalPayment)) - r.recordConfigUsage(ctx, runtimeConfig, &refundOrder) + orderID, originalPayment.OrderNo, refundOrder.ID, bizType, runtimeConfig.Provider, refundAmountCent, refundOrder.ThirdOrderID, originalPayment.ThirdOrderID, refundOriginProviderOrderID(originalPayment)) + r.recordConfigUsage(ctx, runtimeConfig, refundOrder) if err := r.markOrderRefunding(ctx, orderID, refundAmountCent); err != nil { log.Printf("[payment] mark order refunding failed order_id=%d err=%v", orderID, err) } if runtimeConfig.Channel == nil { - _ = r.markRefundFailed(ctx, refundOrder.ID, orderID, refundAmountCent, map[string]string{"error": "payment channel unavailable"}) + if err := r.markRefundFailed(ctx, refundOrder.ID, orderID, refundAmountCent, map[string]string{"error": "payment channel unavailable"}); err != nil { + log.Printf("[payment] mark refund failed status failed order_id=%d payment_id=%d err=%v", orderID, refundOrder.ID, err) + } return nil, ErrPaymentUnavailable } resp, err := runtimeConfig.Channel.CreateRefund(ctx, channelCreateRefundRequest{ ThirdOrderID: originalPayment.ThirdOrderID, ProviderOrderID: refundOriginProviderOrderID(originalPayment), - MerchantRefundID: merchantRefundID, + MerchantRefundID: refundOrder.ThirdOrderID, RefundAmountCent: refundAmountCent, NotifyURL: runtimeConfig.NotifyURL, Attach: originalPayment.OrderNo, Remark: remark, }) if err != nil { - _ = r.markRefundFailed(ctx, refundOrder.ID, orderID, refundAmountCent, map[string]string{"error": err.Error()}) + if markErr := r.markRefundFailed(ctx, refundOrder.ID, orderID, refundAmountCent, map[string]string{"error": err.Error()}); markErr != nil { + log.Printf("[payment] mark refund failed status failed order_id=%d payment_id=%d err=%v", orderID, refundOrder.ID, markErr) + } log.Printf("[payment] refund request failed order_id=%d payment_id=%d biz_type=%s provider=%s amount_cent=%d err=%v", orderID, refundOrder.ID, bizType, runtimeConfig.Provider, refundAmountCent, err) return nil, err } if !resp.OK { - _ = r.markRefundFailed(ctx, refundOrder.ID, orderID, refundAmountCent, resp.Raw) + if markErr := r.markRefundFailed(ctx, refundOrder.ID, orderID, refundAmountCent, resp.Raw); markErr != nil { + log.Printf("[payment] mark refund rejected status failed order_id=%d payment_id=%d err=%v", orderID, refundOrder.ID, markErr) + } log.Printf("[payment] refund rejected order_id=%d payment_id=%d biz_type=%s provider=%s amount_cent=%d code=%s message=%s", orderID, refundOrder.ID, bizType, runtimeConfig.Provider, refundAmountCent, firstNonEmpty(resp.Raw["code"], resp.Raw["resp_code"], resp.Raw["result_code"]), resp.ErrorMessage) return nil, ErrPaymentUnavailable } - refundStatus := "refunding" - var paidAt *time.Time - if resp.Status == "refunded" { - refundStatus = "refunded" - now := time.Now() - paidAt = &now - } else if resp.Status == "failed" { - refundStatus = "failed" - } - if err := r.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", refundOrder.ID).Updates(map[string]any{ - "status": refundStatus, - "provider_order_id": resp.ProviderRefundID, - "raw_request": jsonMap(resp.RawRequest), - "raw_response": jsonMap(withRawSource(resp.Raw, channelSourceCreate)), - "paid_at": paidAt, - }).Error; err != nil { + if err := r.applyRefundChannelStatus(ctx, refundOrder, refundChannelStatusUpdate{ + Status: resp.Status, + ProviderRefundID: resp.ProviderRefundID, + RawRequest: resp.RawRequest, + Raw: resp.Raw, + Source: channelSourceCreate, + }); err != nil { return nil, err } - - if refundStatus == "refunded" { - _ = r.updateOrderRefundStatus(ctx, orderID, refundAmountCent) - refundOrder.PaidAt = paidAt - } else if refundStatus == "failed" { - _ = r.markOrderRefundFailed(ctx, orderID, refundAmountCent) - } else { - _ = r.markOrderRefunding(ctx, orderID, refundAmountCent) - } - refundOrder.Status = refundStatus - refundOrder.ProviderOrderID = resp.ProviderRefundID log.Printf("[payment] refund result order_id=%d payment_id=%d biz_type=%s provider=%s amount_cent=%d status=%s provider_refund_id=%s", - orderID, refundOrder.ID, bizType, runtimeConfig.Provider, refundAmountCent, refundStatus, resp.ProviderRefundID) + orderID, refundOrder.ID, bizType, runtimeConfig.Provider, refundAmountCent, resp.Status, resp.ProviderRefundID) - dto := toRefundDTO(refundOrder) + latest, err := r.findPaymentByID(ctx, refundOrder.ID) + if err != nil { + return nil, err + } + dto := toRefundDTO(*latest) return &dto, nil } func (r *Repository) QueryRefundStatus(ctx context.Context, orderID uint64) (*RefundDTO, error) { @@ -155,49 +125,210 @@ func (r *Repository) QueryRefundStatus(ctx context.Context, orderID uint64) (*Re } return nil, err } - runtimeConfig, err := r.runtimeConfigForPayment(ctx, &payment) + latest, err := r.syncRefundPayment(ctx, &payment, payment.Status == "failed") + if err != nil { + return nil, err + } + dto := toRefundDTO(*latest) + return &dto, nil +} +func (r *Repository) SyncRefundStatusByPaymentID(ctx context.Context, paymentID uint64) (*RefundDTO, error) { + payment, err := r.findPaymentByID(ctx, paymentID) + if err != nil { + if err == gorm.ErrRecordNotFound { + return nil, ErrPaymentNotFound + } + return nil, err + } + if !isRefundBizType(payment.BizType) { + return nil, ErrPaymentNotFound + } + latest, err := r.syncRefundPayment(ctx, payment, true) + if err != nil { + return nil, err + } + dto := toRefundDTO(*latest) + return &dto, nil +} +func (r *Repository) findOriginalPayment(ctx context.Context, orderID uint64) (model.PaymentOrder, error) { + var originalPayment model.PaymentOrder + if err := r.db.WithContext(ctx).Where("order_id = ? AND status = 'paid' AND biz_type = 'order_pay'", orderID).Order("id DESC").First(&originalPayment).Error; err != nil { + if err == gorm.ErrRecordNotFound { + return originalPayment, ErrPaymentNotFound + } + return originalPayment, err + } + return originalPayment, nil +} +func (r *Repository) prepareRefundOrder(ctx context.Context, originalPayment model.PaymentOrder, runtimeConfig runtimePaymentConfig, refundAmountCent int64, bizType string) (*model.PaymentOrder, bool, error) { + var paymentID uint64 + existing := false + err := r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var order model.RentalOrder + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&order, originalPayment.OrderID).Error; err != nil { + return err + } + var existingRefund model.PaymentOrder + err := tx.Clauses(clause.Locking{Strength: "UPDATE"}). + Where("order_id = ? AND biz_type = ?", originalPayment.OrderID, bizType). + Order("id DESC"). + First(&existingRefund).Error + if err == nil { + paymentID = existingRefund.ID + existing = true + return nil + } + if err != gorm.ErrRecordNotFound { + return err + } + paymentNo, err := newPaymentNo() + if err != nil { + return err + } + merchantRefundID := "REF" + paymentNo[3:] + refundOrder := model.PaymentOrder{ + PaymentNo: paymentNo, + OrderID: originalPayment.OrderID, + OrderNo: originalPayment.OrderNo, + UserID: originalPayment.UserID, + Provider: runtimeConfig.Provider, + MerchantID: runtimeConfig.MerchantID, + ThirdOrderID: merchantRefundID, + ProviderOrderID: "", + PayWay: originalPayment.PayWay, + JSPayFlag: originalPayment.JSPayFlag, + AmountCent: refundAmountCent, + BizType: bizType, + Status: "refunding", + } + if err := tx.Create(&refundOrder).Error; err != nil { + return err + } + paymentID = refundOrder.ID + return nil + }) + if err != nil { + return nil, false, err + } + payment, err := r.findPaymentByID(ctx, paymentID) + if err != nil { + return nil, false, err + } + return payment, existing, nil +} +func (r *Repository) syncRefundPayment(ctx context.Context, payment *model.PaymentOrder, queryTerminal bool) (*model.PaymentOrder, error) { + if payment.Status == "refunded" && !queryTerminal { + if err := r.updateOrderRefundStatus(ctx, payment.OrderID, payment.AmountCent); err != nil { + return nil, err + } + return payment, nil + } + if payment.Status == "failed" && !queryTerminal { + if err := r.markOrderRefundFailed(ctx, payment.OrderID, payment.AmountCent); err != nil { + return nil, err + } + return payment, nil + } + runtimeConfig, err := r.runtimeConfigForPayment(ctx, payment) if err != nil { return nil, ErrPaymentUnavailable } - if payment.Status == "refunded" || payment.Status == "failed" || runtimeConfig.isMockMode() { - dto := toRefundDTO(payment) - return &dto, nil + if runtimeConfig.isMockMode() { + if payment.Status == "refunded" { + if err := r.updateOrderRefundStatus(ctx, payment.OrderID, payment.AmountCent); err != nil { + return nil, err + } + } + return payment, nil } if runtimeConfig.Channel == nil { return nil, ErrPaymentUnavailable } - resp, err := runtimeConfig.Channel.QueryRefund(ctx, channelQueryRefundRequest{ - ThirdOrderID: payment.ThirdOrderID, - MerchantRefundID: payment.ThirdOrderID, - ProviderRefundID: payment.ProviderOrderID, - }) + originalPayment, err := r.findOriginalPayment(ctx, payment.OrderID) if err != nil { return nil, err } - if resp.Status == "refunded" { - now := time.Now() - if err := r.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", payment.ID).Updates(map[string]any{ - "status": "refunded", - "paid_at": now, - "raw_response": jsonMap(withRawSource(resp.Raw, channelSourceQuery)), - }).Error; err != nil { - return nil, err + resp, err := runtimeConfig.Channel.QueryRefund(ctx, refundQueryRequest(*payment, originalPayment)) + if err != nil { + return nil, err + } + if !resp.OK { + return nil, ErrPaymentUnavailable + } + if err := r.applyRefundChannelStatus(ctx, payment, refundChannelStatusUpdate{ + Status: resp.Status, + ProviderRefundID: resp.ProviderRefundID, + RefundTime: resp.RefundTime, + Raw: resp.Raw, + Source: channelSourceQuery, + }); err != nil { + return nil, err + } + return r.findPaymentByID(ctx, payment.ID) +} + +type refundChannelStatusUpdate struct { + Status string + ProviderRefundID string + RefundTime string + RawRequest map[string]string + Raw map[string]string + Source string +} + +func (r *Repository) applyRefundChannelStatus(ctx context.Context, payment *model.PaymentOrder, update refundChannelStatusUpdate) error { + status := update.Status + if status != "refunded" && status != "failed" { + status = "refunding" + } + updates := map[string]any{ + "status": status, + "raw_response": jsonMap(withRawSource(update.Raw, update.Source)), + } + if update.ProviderRefundID != "" { + updates["provider_order_id"] = update.ProviderRefundID + } + if update.RawRequest != nil { + updates["raw_request"] = jsonMap(update.RawRequest) + } + if status == "refunded" { + paidAt := parseChannelTime(update.RefundTime) + if paidAt == nil { + now := time.Now() + paidAt = &now } - payment.Status = "refunded" - payment.PaidAt = &now - _ = r.updateOrderRefundStatus(ctx, orderID, payment.AmountCent) - } else if resp.Status == "failed" { - if err := r.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", payment.ID).Updates(map[string]any{ - "status": "failed", - "raw_response": jsonMap(withRawSource(resp.Raw, channelSourceQuery)), - }).Error; err != nil { - return nil, err + updates["paid_at"] = paidAt + } + if update.Source == channelSourceNotify { + updates["notified_at"] = time.Now() + } + if err := r.db.WithContext(ctx).Model(&model.PaymentOrder{}).Where("id = ?", payment.ID).Updates(updates).Error; err != nil { + return err + } + switch status { + case "refunded": + return r.updateOrderRefundStatus(ctx, payment.OrderID, payment.AmountCent) + case "failed": + return r.markOrderRefundFailed(ctx, payment.OrderID, payment.AmountCent) + default: + return r.markOrderRefunding(ctx, payment.OrderID, payment.AmountCent) + } +} +func isRefundBizType(bizType string) bool { + for _, item := range refundBizTypes { + if item == bizType { + return true } - payment.Status = "failed" - _ = r.markOrderRefundFailed(ctx, orderID, payment.AmountCent) + } + return false +} +func refundQueryRequest(payment model.PaymentOrder, originalPayment model.PaymentOrder) channelQueryRefundRequest { + return channelQueryRefundRequest{ + ThirdOrderID: originalPayment.ThirdOrderID, + ProviderOrderID: refundOriginProviderOrderID(originalPayment), + MerchantRefundID: payment.ThirdOrderID, + ProviderRefundID: payment.ProviderOrderID, } - dto := toRefundDTO(payment) - return &dto, nil } func (r *Repository) updateOrderRefundStatus(ctx context.Context, orderID uint64, refundAmountCent int64) error { now := time.Now() diff --git a/backend/internal/modules/payment/repository.go b/backend/internal/modules/payment/repository.go index da20095..c02cdb2 100644 --- a/backend/internal/modules/payment/repository.go +++ b/backend/internal/modules/payment/repository.go @@ -44,6 +44,12 @@ var refundBizTypes = []string{ "arbitration_refund", } +func RefundBizTypes() []string { + out := make([]string, len(refundBizTypes)) + copy(out, refundBizTypes) + return out +} + func NewRepository(db *gorm.DB, configRepo *paymentconfig.Repository, orderRepo *order.Repository) *Repository { return &Repository{ db: db, diff --git a/backend/internal/modules/payment/repository_integration_test.go b/backend/internal/modules/payment/repository_integration_test.go index e9d8a07..2f6f6f7 100644 --- a/backend/internal/modules/payment/repository_integration_test.go +++ b/backend/internal/modules/payment/repository_integration_test.go @@ -282,3 +282,67 @@ func TestNotifyResultHasRequiredFields(t *testing.T) { t.Fatal("NotifyResult.Message should not be empty") } } + +func TestStartRefundDoesNotCreateNewOrderWhenFailedRefundExists(t *testing.T) { + db := setupPaymentTestDB(t) + repo := NewRepository(db, nil, nil) + order := model.RentalOrder{ + ID: 1001, + OrderNo: "ORD202606140001", + RenterID: 11, + RentAmountCent: 800, + DepositAmountCent: 200, + RefundStatus: "failed", + RefundAmountCent: 1000, + Status: "closed", + } + if err := db.Create(&order).Error; err != nil { + t.Fatalf("create order failed: %v", err) + } + original := model.PaymentOrder{ + PaymentNo: "PAY202606140001", + OrderID: order.ID, + OrderNo: order.OrderNo, + UserID: order.RenterID, + Provider: "mock", + ThirdOrderID: "PAY202606140001", + ProviderOrderID: "MOCKPAY202606140001", + AmountCent: 1000, + BizType: "order_pay", + Status: "paid", + } + if err := db.Create(&original).Error; err != nil { + t.Fatalf("create original payment failed: %v", err) + } + existingRefund := model.PaymentOrder{ + PaymentNo: "PAY202606140002", + OrderID: order.ID, + OrderNo: order.OrderNo, + UserID: order.RenterID, + Provider: "mock", + ThirdOrderID: "REF202606140002", + AmountCent: 1000, + BizType: "admin_refund", + Status: "failed", + } + if err := db.Create(&existingRefund).Error; err != nil { + t.Fatalf("create existing refund failed: %v", err) + } + + dto, err := repo.StartRefund(t.Context(), order.ID, 1000, "admin_refund", "后台人工退款") + if err != nil { + t.Fatalf("StartRefund() error = %v", err) + } + if dto.ID != existingRefund.ID { + t.Fatalf("StartRefund() returned payment id %d, want existing %d", dto.ID, existingRefund.ID) + } + var count int64 + if err := db.Model(&model.PaymentOrder{}). + Where("order_id = ? AND biz_type = ?", order.ID, "admin_refund"). + Count(&count).Error; err != nil { + t.Fatalf("count refund orders failed: %v", err) + } + if count != 1 { + t.Fatalf("refund order count = %d, want 1", count) + } +} diff --git a/backend/internal/modules/payment/repository_test.go b/backend/internal/modules/payment/repository_test.go index 62ae256..dfa72df 100644 --- a/backend/internal/modules/payment/repository_test.go +++ b/backend/internal/modules/payment/repository_test.go @@ -157,3 +157,25 @@ func TestRefundOriginProviderOrderIDFallsBackToProviderOrderID(t *testing.T) { t.Fatalf("refundOriginProviderOrderID() = %q, want fallback %q", got, payment.ProviderOrderID) } } + +func TestRefundQueryRequestUsesOriginalPaymentAndMerchantRefundID(t *testing.T) { + original := model.PaymentOrder{ + ThirdOrderID: "PAY202606140001", + ProviderOrderID: "PROVIDER-PAY-ID", + } + refund := model.PaymentOrder{ + ThirdOrderID: "REF202606140001", + ProviderOrderID: "PROVIDER-REFUND-ID", + } + + req := refundQueryRequest(refund, original) + if req.ThirdOrderID != original.ThirdOrderID { + t.Fatalf("ThirdOrderID = %q, want original payment third order id %q", req.ThirdOrderID, original.ThirdOrderID) + } + if req.MerchantRefundID != refund.ThirdOrderID { + t.Fatalf("MerchantRefundID = %q, want refund third order id %q", req.MerchantRefundID, refund.ThirdOrderID) + } + if req.ProviderRefundID != refund.ProviderOrderID { + t.Fatalf("ProviderRefundID = %q, want refund provider order id %q", req.ProviderRefundID, refund.ProviderOrderID) + } +} diff --git a/backend/migrations/000003_refund_retry_index.sql b/backend/migrations/000003_refund_retry_index.sql new file mode 100644 index 0000000..4f9dff1 --- /dev/null +++ b/backend/migrations/000003_refund_retry_index.sql @@ -0,0 +1,9 @@ +-- 退款补偿重试字段:限制失败补偿的自动重试次数,并支持退避。 +ALTER TABLE payment_orders + ADD COLUMN retry_count INT NOT NULL DEFAULT 0 COMMENT '自动补偿重试次数' AFTER raw_response, + ADD COLUMN last_retry_at DATETIME NULL COMMENT '最近自动补偿重试时间' AFTER retry_count, + ADD COLUMN next_retry_at DATETIME NULL COMMENT '下次允许自动补偿重试时间' AFTER last_retry_at; + +-- 退款补偿扫描索引:按业务类型、状态、下次重试时间批量捞取待同步退款单。 +CREATE INDEX idx_payment_orders_refund_retry + ON payment_orders (biz_type, status, next_retry_at, updated_at, id);