feat(pay-v2): P8 Task6 拒付 chargeback 记录 + chargeback.received

This commit is contained in:
wangjia
2026-07-10 18:52:13 +08:00
parent 02b2fcfa41
commit b0714ca758
19 changed files with 502 additions and 22 deletions
+171
View File
@@ -0,0 +1,171 @@
package gateway_test
import (
"context"
"encoding/json"
"testing"
"github.com/wangjia/pay/config"
"github.com/wangjia/pay/internal/accounts"
"github.com/wangjia/pay/internal/gateway"
"github.com/wangjia/pay/internal/model"
"github.com/wangjia/pay/internal/provider"
"github.com/wangjia/pay/internal/store"
)
// newChargebackGateway 装配一套独立的 gateway(渠道 "substripe",复用 subscription_test.go
// 的 fakeSubProvider——它把测试注入的 JSON 原样反序列化成 provider.PaidEvent,包括
// Kind/DisputeRef/OutTradeNo 等 P8 新增字段,fake.Provider 的精简版协议做不到这点),额外
// 暴露 *store.ChargebackStore 供断言落库情况(P8 Task6 专用,不复用 newGateway/newSubGateway
// 避免改动其多处既有调用签名)。
func newChargebackGateway(t *testing.T) (*gateway.Gateway, *store.OrderStore, *store.ChargebackStore, *spyEnqueuer) {
t.Helper()
db := model.OpenTestDB(t)
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
preg := provider.NewRegistry()
preg.Register(&fakeSubProvider{sessionRef: "cs_cb_1"})
areg := accounts.New([]config.AccountConfig{
{AccountID: "cb-a1", Channel: "substripe", Region: "global", Enabled: true, Weight: 1},
})
picker := accounts.NewRouter(areg, nil, nil)
spy := &spyEnqueuer{}
g := gateway.New(orders, refunds, preg, picker, stubSubResolver{}, spy, "global", subs, chargebacks)
return g, orders, chargebacks, spy
}
func seedPaidOrder(t *testing.T, orders *store.OrderStore, no string) {
t.Helper()
if err := orders.CreateOrder(&model.OrderV2{
OutTradeNo: no, BizSystem: "pangolin", BizRef: "u-1", BizCode: "pro_month",
AmountMinor: 2999, Currency: "USD", Status: model.OrderPaidV2,
}); err != nil {
t.Fatalf("seed paid order: %v", err)
}
}
// TestRecordChargebackHappyPath 覆盖 brief Step1 ①②:charge.dispute.created → 落
// Chargeback 一行 + 原 order Disputed=true + webhook spy 收 chargeback.received,且不
// 自动改订单状态机(仍是 paid,只是 Disputed 打标)。
func TestRecordChargebackHappyPath(t *testing.T) {
g, orders, chargebacks, spy := newChargebackGateway(t)
seedPaidOrder(t, orders, "PAY-1")
raw, err := json.Marshal(provider.PaidEvent{
Kind: provider.EventChargeback, DisputeRef: "dp_1", OutTradeNo: "PAY-1",
ProviderPaymentRef: "pi_1", PaidAmountMinor: 2999, PaidCurrency: "USD",
Reason: "fraudulent", Status: provider.PaidFailed,
})
if err != nil {
t.Fatalf("marshal: %v", err)
}
result, err := g.HandleCallback(context.Background(), "substripe", provider.CallbackInput{Raw: raw})
if err != nil || result != gateway.SettleProcessed {
t.Fatalf("HandleCallback = %v, %v, want SettleProcessed", result, err)
}
o, err := orders.GetOrder("PAY-1")
if err != nil {
t.Fatalf("GetOrder: %v", err)
}
if !o.Disputed {
t.Fatalf("order.Disputed = false, want true")
}
if o.Status != model.OrderPaidV2 {
t.Fatalf("order.Status = %s, want unchanged paid(打标不改状态机)", o.Status)
}
_ = chargebacks // 幂等落库由下面的重投用例断言(created=false)
if len(spy.calls) != 1 {
t.Fatalf("webhook calls = %d, want 1: %+v", len(spy.calls), spy.calls)
}
c := spy.calls[0]
if c["event_type"] != gateway.EvtChargebackReceived || c["out_trade_no"] != "PAY-1" || c["dispute_ref"] != "dp_1" {
t.Fatalf("webhook payload = %+v", c)
}
}
// TestRecordChargebackDuplicateNotDoubleRecordedOrEnqueued 拒付重投(Stripe 常见重投场景)→
// Chargeback 表不双记(ON CONFLICT dispute_ref)、webhook 不双发(outbox 唯一键)。
func TestRecordChargebackDuplicateNotDoubleRecordedOrEnqueued(t *testing.T) {
g, orders, _, spy := newChargebackGateway(t)
seedPaidOrder(t, orders, "PAY-2")
raw, err := json.Marshal(provider.PaidEvent{
Kind: provider.EventChargeback, DisputeRef: "dp_2", OutTradeNo: "PAY-2",
ProviderPaymentRef: "pi_2", PaidAmountMinor: 1999, PaidCurrency: "USD",
Reason: "duplicate", Status: provider.PaidFailed,
})
if err != nil {
t.Fatalf("marshal: %v", err)
}
if _, err := g.HandleCallback(context.Background(), "substripe", provider.CallbackInput{Raw: raw}); err != nil {
t.Fatalf("HandleCallback#1: %v", err)
}
result2, err := g.HandleCallback(context.Background(), "substripe", provider.CallbackInput{Raw: raw})
if err != nil {
t.Fatalf("HandleCallback#2: %v", err)
}
if result2 != gateway.SettleDuplicate {
t.Fatalf("replay result = %v, want duplicate", result2)
}
if len(spy.calls) != 1 {
t.Fatalf("webhook calls after replay = %d, want still 1: %+v", len(spy.calls), spy.calls)
}
}
// TestRecordChargebackEmptyOutTradeNoNotEnqueued 覆盖 brief Step1 ③:out_trade_no 空
// (订阅拒付,PI 无 metadata)→ 仍落 Chargeback 记录,但不入队业务 webhook(无法定位业务单)。
func TestRecordChargebackEmptyOutTradeNoNotEnqueued(t *testing.T) {
g, _, chargebacks, spy := newChargebackGateway(t)
raw, err := json.Marshal(provider.PaidEvent{
Kind: provider.EventChargeback, DisputeRef: "dp_sub_1", OutTradeNo: "",
ProviderPaymentRef: "pi_sub_1", PaidAmountMinor: 999, PaidCurrency: "USD",
Reason: "fraudulent", Status: provider.PaidFailed,
})
if err != nil {
t.Fatalf("marshal: %v", err)
}
result, err := g.HandleCallback(context.Background(), "substripe", provider.CallbackInput{Raw: raw})
if err != nil || result != gateway.SettleProcessed {
t.Fatalf("HandleCallback = %v, %v, want SettleProcessed(已记录未转发)", result, err)
}
if len(spy.calls) != 0 {
t.Fatalf("webhook calls = %d, want 0(无法定位业务单不转发): %+v", len(spy.calls), spy.calls)
}
// 重投同一空 out_trade_no dispute → Chargeback 仍不双记(created 幂等),同样不入队。
again, err := g.HandleCallback(context.Background(), "substripe", provider.CallbackInput{Raw: raw})
if err != nil || again != gateway.SettleDuplicate {
t.Fatalf("replay = %v, %v, want duplicate", again, err)
}
if len(spy.calls) != 0 {
t.Fatalf("webhook calls after replay = %d, want still 0", len(spy.calls))
}
_ = chargebacks
}
// TestRecordChargebackUnknownOrderNotBlocking out_trade_no 非空但查单失败(极端场景,如脏
// 数据/竞态)→ 已记录 Chargeback,定位失败不阻断、不 panic,同样不转发。
func TestRecordChargebackUnknownOrderNotBlocking(t *testing.T) {
g, _, _, spy := newChargebackGateway(t)
raw, err := json.Marshal(provider.PaidEvent{
Kind: provider.EventChargeback, DisputeRef: "dp_unknown", OutTradeNo: "NOPE",
ProviderPaymentRef: "pi_unknown", PaidAmountMinor: 500, PaidCurrency: "USD",
Reason: "fraudulent", Status: provider.PaidFailed,
})
if err != nil {
t.Fatalf("marshal: %v", err)
}
result, err := g.HandleCallback(context.Background(), "substripe", provider.CallbackInput{Raw: raw})
if err != nil || result != gateway.SettleProcessed {
t.Fatalf("HandleCallback = %v, %v, want SettleProcessed(已记录,查单失败不阻断)", result, err)
}
if len(spy.calls) != 0 {
t.Fatalf("webhook calls = %d, want 0", len(spy.calls))
}
}
+2 -1
View File
@@ -46,6 +46,7 @@ func TestE2ECryptoQuerySettles(t *testing.T) {
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
acctReg := accounts.New([]config.AccountConfig{
{AccountID: "e2e-1", Channel: "crypto", Enabled: true, Region: "global", CredentialEnvPrefix: "e2e"},
})
@@ -57,7 +58,7 @@ func TestE2ECryptoQuerySettles(t *testing.T) {
picker := accounts.NewRouter(acctReg, nil, nil)
spy := &spyEnqueuer{}
g := gateway.New(orders, refunds, preg, picker, cryptoResolver{}, spy, "global", subs)
g := gateway.New(orders, refunds, preg, picker, cryptoResolver{}, spy, "global", subs, chargebacks)
// 下单 → 从 session payload 拿到期望链上金额(base+唯一尾数),喂给假 TronGrid。
res, err := g.CreateOrder(context.Background(), gateway.CreateOrderInput{
+12 -10
View File
@@ -39,20 +39,22 @@ type WebhookEnqueuer interface {
}
type Gateway struct {
orders *store.OrderStore
refunds *store.RefundStore
providers *provider.Registry
picker accounts.Picker
products ProductResolver
webhook WebhookEnqueuer
region string
subs *store.SubscriptionStore
orders *store.OrderStore
refunds *store.RefundStore
providers *provider.Registry
picker accounts.Picker
products ProductResolver
webhook WebhookEnqueuer
region string
subs *store.SubscriptionStore
chargebacks *store.ChargebackStore
}
func New(orders *store.OrderStore, refunds *store.RefundStore, providers *provider.Registry, picker accounts.Picker,
products ProductResolver, webhook WebhookEnqueuer, region string, subs *store.SubscriptionStore) *Gateway {
products ProductResolver, webhook WebhookEnqueuer, region string, subs *store.SubscriptionStore,
chargebacks *store.ChargebackStore) *Gateway {
return &Gateway{orders: orders, refunds: refunds, providers: providers, picker: picker,
products: products, webhook: webhook, region: region, subs: subs}
products: products, webhook: webhook, region: region, subs: subs, chargebacks: chargebacks}
}
type CreateOrderInput struct {
+2 -1
View File
@@ -73,6 +73,7 @@ func newGateway(t *testing.T) (*gateway.Gateway, *fake.Provider, *spyEnqueuer, *
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
preg := provider.NewRegistry()
fp := fake.New()
preg.Register(fp)
@@ -83,7 +84,7 @@ func newGateway(t *testing.T) (*gateway.Gateway, *fake.Provider, *spyEnqueuer, *
})
picker := accounts.NewRouter(areg, nil, nil) // 默认 round_robin
spy := &spyEnqueuer{}
g := gateway.New(orders, refunds, preg, picker, stubResolver{}, spy, "global", subs)
g := gateway.New(orders, refunds, preg, picker, stubResolver{}, spy, "global", subs, chargebacks)
return g, fp, spy, orders
}
+52 -2
View File
@@ -3,7 +3,6 @@ package gateway
import (
"context"
"errors"
"fmt"
"log"
"time"
@@ -207,8 +206,59 @@ func (g *Gateway) enqueueRenewed(sub *model.Subscription, renewalNo string, ev *
})
}
// recordChargeback 处理入站 charge.dispute.created(P8 Task6,设计 §6「钱到账不可逆」的
// 例外形态)。只 Stripe(卡)有此语义;alipay/微信本轮无拒付流,crypto 收款永无
// chargeback——三者均不产出 EventChargeback,本函数只会被 stripe adapter 触发。
//
// 决策记录:不自动回收权益(§6,业务方裁量),只做三件事——①落 Chargeback(幂等 by
// dispute_ref,DisputeRef 重投 no-op,不双记)②给原订单打 Disputed 标(不改状态机)
// ③入队 chargeback.received 给业务方自行冲正。out_trade_no 解析失败(订阅拒付/查单失败)
// 时仍落 Chargeback 留痕 + log 告警,但不阻断、不转发(无法定位业务单)。
//
// 入队不按 created 分叉(与 settleRenewal/finalizeCanceled/markSubscriptionPastDue 同型的
// P8 反纪律修法):Chargeback.Create 与 Enqueue 是两次独立写,不在同一事务——若首次 Create
// 成功但 Enqueue 瞬时失败,调用方(HandleCallback)拿到 err 后 Stripe 会重投同一
// charge.dispute.created,此时 created 必为 false;若像"created=false→直接 return
// SettleDuplicate"那样跳过下面的打标+入队,chargeback.received 通知永久丢失、Disputed
// 标也永远打不上。改为无论 created 与否都走完打标+入队,outbox 唯一键 ON CONFLICT DO
// NOTHING + MarkDisputed 条件 UPDATE 天然双重幂等——重投即自愈,不会双记/双发。
func (g *Gateway) recordChargeback(ctx context.Context, method string, ev *provider.PaidEvent) (SettleResult, error) {
return SettleFailed, fmt.Errorf("not implemented: %s", ev.Kind)
created, err := g.chargebacks.Create(&model.Chargeback{
DisputeRef: ev.DisputeRef, OutTradeNo: ev.OutTradeNo, Channel: method,
ProviderPaymentRef: ev.ProviderPaymentRef, AmountMinor: ev.PaidAmountMinor,
Currency: ev.PaidCurrency, Reason: ev.Reason, Status: "received",
})
if err != nil {
return SettleFailed, err
}
result := SettleProcessed
if !created {
result = SettleDuplicate // 拒付重投:Chargeback 已记录过,但仍需补齐下面的打标/入队(见函数注释)
}
if ev.OutTradeNo == "" {
log.Printf("[chargeback] dispute=%s 无法定位业务单(订阅/无 metadata),已记录未转发", ev.DisputeRef)
return result, nil
}
o, err := g.orders.GetOrder(ev.OutTradeNo)
if err != nil {
log.Printf("[chargeback] dispute=%s out_trade_no=%s 查单失败: %v", ev.DisputeRef, ev.OutTradeNo, err)
return result, nil // 已记录 chargeback;定位失败不阻断
}
if _, err := g.orders.MarkDisputed(o.OutTradeNo); err != nil { // 打标不改状态机;条件 UPDATE 幂等
return SettleFailed, err
}
if o.BizSystem == "" {
return result, nil
}
if err := g.webhook.Enqueue(o.OutTradeNo, o.BizSystem, EvtChargebackReceived, "", map[string]any{
"event_type": EvtChargebackReceived, "out_trade_no": o.OutTradeNo, "dispute_ref": ev.DisputeRef,
"biz_system": o.BizSystem, "biz_ref": o.BizRef, "product_biz_code": o.BizCode,
"amount_minor": ev.PaidAmountMinor, "currency": ev.PaidCurrency, "reason": ev.Reason,
"received_at": time.Now().Format(time.RFC3339),
}); err != nil {
return SettleFailed, err
}
return result, nil
}
// SyncPendingAttempts polls every pending attempt via its Provider.Query and
+2 -1
View File
@@ -139,6 +139,7 @@ func TestSettleTransientReadErrorIsFailed(t *testing.T) {
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
preg := provider.NewRegistry()
fp := fake.New()
preg.Register(fp)
@@ -147,7 +148,7 @@ func TestSettleTransientReadErrorIsFailed(t *testing.T) {
})
picker := accounts.NewRouter(areg, nil, nil)
spy := &spyEnqueuer{}
g := gateway.New(orders, refunds, preg, picker, stubResolver{}, spy, "global", subs)
g := gateway.New(orders, refunds, preg, picker, stubResolver{}, spy, "global", subs, chargebacks)
ctx := context.Background()
g.CreateOrder(ctx, gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-1"})
+2 -1
View File
@@ -74,6 +74,7 @@ func newSubGateway(t *testing.T) (*gateway.Gateway, *fakeSubProvider, *spyEnqueu
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
preg := provider.NewRegistry()
fp := &fakeSubProvider{sessionRef: "cs_test_sess1"}
preg.Register(fp)
@@ -82,7 +83,7 @@ func newSubGateway(t *testing.T) (*gateway.Gateway, *fakeSubProvider, *spyEnqueu
})
picker := accounts.NewRouter(areg, nil, nil)
spy := &spyEnqueuer{}
g := gateway.New(orders, refunds, preg, picker, stubSubResolver{}, spy, "global", subs)
g := gateway.New(orders, refunds, preg, picker, stubSubResolver{}, spy, "global", subs, chargebacks)
return g, fp, spy, orders, subs
}