0ae7769abe
- gateway: TestSettleConcurrentSameOrder 真并发钉住 MarkAttemptPaid 行锁不变量 (文件型 sqlite,规避 in-memory cache=shared 的 SQLITE_LOCKED_SHAREDCACHE)。 - gateway: e2e_fullchain_test.go 补下单→回调→settle→webhook 实际 HTTP 投递→ delivered 整链(httptest server + 真实 store.WebhookStore/webhook.Notifier)。 - reconcile: main.go 装配抽到 reconcile.Assemble(+Runner.TaskNames 访问器), 补 assembly_test.go 钉住 7 个后台任务全部注册 + crypto 缺渠道/interval=0 分支。 - alipay: 补验签健壮性(未知多余字段/字段乱序/中文unicode/空值字段),意外定位 vendor Encoder 对空值字段的真实语义与官方文档描述不同,已记录在测试注释。 - money: 补 TestZeroDecimalCurrencyJPYNotSupported,记录 JPY/KRW 当前未注册进 exponents(接入会先踩 ErrUnknownCurrency,不是乘除法坑)。
318 lines
12 KiB
Go
318 lines
12 KiB
Go
package gateway_test
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"sync"
|
|
"testing"
|
|
|
|
"github.com/glebarez/sqlite"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/logger"
|
|
|
|
"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/provider/fake"
|
|
"github.com/wangjia/pay/internal/store"
|
|
"github.com/wangjia/pay/internal/webhook"
|
|
)
|
|
|
|
func attemptRef(t *testing.T, orders interface {
|
|
ListAttemptsByStatus(model.AttemptStatus, int) ([]model.Attempt, error)
|
|
}) string {
|
|
t.Helper()
|
|
atts, _ := orders.ListAttemptsByStatus(model.AttemptPending, 10)
|
|
if len(atts) == 0 {
|
|
t.Fatalf("无 pending 尝试")
|
|
}
|
|
return atts[0].ProviderRef
|
|
}
|
|
|
|
func TestSettleHappyIdempotentAndWebhook(t *testing.T) {
|
|
g, _, spy, orders := newGateway(t)
|
|
ctx := context.Background()
|
|
res, _ := g.CreateOrder(ctx, gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-1"})
|
|
ref := attemptRef(t, orders)
|
|
|
|
ev := &provider.PaidEvent{ProviderRef: ref, Status: provider.PaidSucceeded, PaidAmountMinor: 29990000, PaidCurrency: "USDT"}
|
|
got, err := g.Settle(ctx, ev)
|
|
if err != nil || got != gateway.SettleProcessed {
|
|
t.Fatalf("settle#1 = %v, %v", got, err)
|
|
}
|
|
// 订单已 paid
|
|
o, _ := orders.GetOrder(res.OrderNo)
|
|
if o.Status != model.OrderPaidV2 {
|
|
t.Fatalf("order 应 paid, got %v", o.Status)
|
|
}
|
|
// webhook 入队一次,payload 带 event_type
|
|
if len(spy.calls) != 1 || spy.calls[0]["event_type"] != "payment.succeeded" || spy.calls[0]["out_trade_no"] != res.OrderNo {
|
|
t.Fatalf("webhook calls = %+v", spy.calls)
|
|
}
|
|
|
|
// 幂等:再 settle → duplicate,不重复入队
|
|
got2, _ := g.Settle(ctx, ev)
|
|
if got2 != gateway.SettleDuplicate || len(spy.calls) != 1 {
|
|
t.Fatalf("settle#2 = %v, calls=%d", got2, len(spy.calls))
|
|
}
|
|
}
|
|
|
|
func TestSettleWebhookCarriesProductBizCode(t *testing.T) {
|
|
g, fp, spy, orders := newGateway(t)
|
|
ctx := context.Background()
|
|
res, _ := g.CreateOrder(ctx, gateway.CreateOrderInput{
|
|
SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-9",
|
|
})
|
|
ref := attemptRef(t, orders)
|
|
fp.SetQueryResult(ref, provider.PaidEvent{
|
|
ProviderRef: ref, Status: provider.PaidSucceeded,
|
|
PaidAmountMinor: 29990000, PaidCurrency: "USDT",
|
|
})
|
|
if _, err := g.SyncPendingAttempts(ctx, 10); err != nil {
|
|
t.Fatalf("sync: %v", err)
|
|
}
|
|
if len(spy.calls) != 1 {
|
|
t.Fatalf("want 1 webhook, got %d", len(spy.calls))
|
|
}
|
|
if spy.calls[0]["product_biz_code"] != "pro_year" {
|
|
t.Fatalf("payload product_biz_code = %v, want pro_year", spy.calls[0]["product_biz_code"])
|
|
}
|
|
_ = res
|
|
}
|
|
|
|
func TestSettleGuards(t *testing.T) {
|
|
g, _, spy, orders := newGateway(t)
|
|
ctx := context.Background()
|
|
g.CreateOrder(ctx, gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-1"})
|
|
ref := attemptRef(t, orders)
|
|
|
|
// 未 succeeded → ignored
|
|
if got, _ := g.Settle(ctx, &provider.PaidEvent{ProviderRef: ref, Status: provider.PaidPending}); got != gateway.SettleIgnored {
|
|
t.Fatalf("pending 应 ignored, got %v", got)
|
|
}
|
|
// 未知 ref → not_found
|
|
if got, _ := g.Settle(ctx, &provider.PaidEvent{ProviderRef: "GHOST", Status: provider.PaidSucceeded, PaidCurrency: "USDT", PaidAmountMinor: 1}); got != gateway.SettleNotFound {
|
|
t.Fatalf("未知 ref 应 not_found, got %v", got)
|
|
}
|
|
// 少付 → amount_mismatch
|
|
if got, err := g.Settle(ctx, &provider.PaidEvent{ProviderRef: ref, Status: provider.PaidSucceeded, PaidCurrency: "USDT", PaidAmountMinor: 1}); got != gateway.SettleAmountMismatch || err == nil {
|
|
t.Fatalf("少付应 amount_mismatch, got %v %v", got, err)
|
|
}
|
|
// 错币种 → amount_mismatch
|
|
if got, _ := g.Settle(ctx, &provider.PaidEvent{ProviderRef: ref, Status: provider.PaidSucceeded, PaidCurrency: "CNY", PaidAmountMinor: 29990000}); got != gateway.SettleAmountMismatch {
|
|
t.Fatalf("错币种应 amount_mismatch, got %v", got)
|
|
}
|
|
if len(spy.calls) != 0 {
|
|
t.Fatalf("守卫失败路径不应入队 webhook, got %d", len(spy.calls))
|
|
}
|
|
}
|
|
|
|
// 资金命脉不变量:outbox 入队失败 → 绝不翻转订单(否则"已付但永不通知")。
|
|
// 渠道拿不到 200 会重投,重投时入队+翻转都幂等,自然恢复。
|
|
func TestSettleEnqueueFailureKeepsOrderPending(t *testing.T) {
|
|
g, _, spy, orders := newGateway(t)
|
|
ctx := context.Background()
|
|
res, _ := g.CreateOrder(ctx, gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-1"})
|
|
ref := attemptRef(t, orders)
|
|
ev := &provider.PaidEvent{ProviderRef: ref, Status: provider.PaidSucceeded, PaidAmountMinor: 29990000, PaidCurrency: "USDT"}
|
|
|
|
spy.failNext = true
|
|
if got, err := g.Settle(ctx, ev); got != gateway.SettleFailed || err == nil {
|
|
t.Fatalf("入队失败应 SettleFailed+err, got %v, %v", got, err)
|
|
}
|
|
o, _ := orders.GetOrder(res.OrderNo)
|
|
if o.Status != model.OrderPendingV2 {
|
|
t.Fatalf("入队失败后订单必须仍 pending, got %v", o.Status)
|
|
}
|
|
|
|
// 渠道重投 → 入队成功 → 翻转
|
|
if got, err := g.Settle(ctx, ev); err != nil || got != gateway.SettleProcessed {
|
|
t.Fatalf("重投应 processed, got %v, %v", got, err)
|
|
}
|
|
if len(spy.calls) != 1 {
|
|
t.Fatalf("恢复后应恰入队 1 次, got %d", len(spy.calls))
|
|
}
|
|
}
|
|
|
|
// 瞬时读库失败(非 store.ErrAttemptNotFound 哨兵)必须归 SettleFailed(可重试),
|
|
// 不能与"查无此单"的终态 SettleNotFound 混淆——否则调用方按结果值决定 ack,
|
|
// 会把已付订单永久丢弃。构造方式:先建好订单/尝试,再直接关掉底层连接,
|
|
// 让 AttemptByProviderRef 打到一个已关闭的 DB 上,产出非哨兵错误。
|
|
func TestSettleTransientReadErrorIsFailed(t *testing.T) {
|
|
db := model.OpenTestDB(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)
|
|
areg := accounts.New([]config.AccountConfig{
|
|
{AccountID: "fake-a1", Channel: "fake", Region: "global", Enabled: true, Weight: 1},
|
|
})
|
|
picker := accounts.NewRouter(areg, nil, nil)
|
|
spy := &spyEnqueuer{}
|
|
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"})
|
|
ref := attemptRef(t, orders)
|
|
|
|
sqlDB, err := db.DB()
|
|
if err != nil {
|
|
t.Fatalf("db.DB(): %v", err)
|
|
}
|
|
if err := sqlDB.Close(); err != nil {
|
|
t.Fatalf("close db: %v", err)
|
|
}
|
|
|
|
ev := &provider.PaidEvent{ProviderRef: ref, Status: provider.PaidSucceeded, PaidAmountMinor: 29990000, PaidCurrency: "USDT"}
|
|
got, err := g.Settle(ctx, ev)
|
|
if got != gateway.SettleFailed || err == nil {
|
|
t.Fatalf("瞬时读库失败应 SettleFailed+err, got %v, %v", got, err)
|
|
}
|
|
}
|
|
|
|
// openFileGatewaySettleDB 开一个 t.TempDir 下的文件型 sqlite(而非 model.OpenTestDB 的
|
|
// in-memory cache=shared)专供并发 Settle 测试:cache=shared 的 in-memory sqlite 下,两个
|
|
// goroutine 各自事务并发写不同表会报 "database table is locked"(SQLITE_LOCKED_SHAREDCACHE,
|
|
// 共享缓存表级锁——不是 busy_timeout 能重试的 SQLITE_BUSY),同 store/refund_test.go
|
|
// openFileGuardedDB 的说明。文件型连接各走独立锁路径,规避此问题;DSN 同样带
|
|
// _txlock=immediate 保留"事务一开始即抢写锁"的守卫语义。
|
|
func openFileGatewaySettleDB(t *testing.T) *gorm.DB {
|
|
t.Helper()
|
|
dsn := fmt.Sprintf("file:%s/settle.db?_txlock=immediate&_pragma=busy_timeout(5000)", t.TempDir())
|
|
db, err := gorm.Open(sqlite.Open(dsn),
|
|
&gorm.Config{Logger: logger.Default.LogMode(logger.Silent), TranslateError: true})
|
|
if err != nil {
|
|
t.Fatalf("open file settle db: %v", err)
|
|
}
|
|
if err := db.AutoMigrate(&model.OrderV2{}, &model.Attempt{}, &model.Account{}, &model.Refund{},
|
|
&model.WebhookDelivery{}, &model.Product{}, &model.ProductPrice{}, &model.OrphanPayment{},
|
|
&model.Subscription{}, &model.Chargeback{}); err != nil {
|
|
t.Fatalf("migrate: %v", err)
|
|
}
|
|
if err := model.UpgradeWebhookDeliveryIndex(db); err != nil {
|
|
t.Fatalf("upgrade webhook_deliveries uq_delivery index: %v", err)
|
|
}
|
|
sqlDB, _ := db.DB()
|
|
t.Cleanup(func() { _ = sqlDB.Close() })
|
|
return db
|
|
}
|
|
|
|
// TestSettleConcurrentSameOrder 钉住 Settle 依赖的行锁不变量(MarkAttemptPaid 用条件
|
|
// UPDATE ... WHERE status=pending 天然把并发翻转串行化,同 store/refund_test.go
|
|
// TestCreateRefundGuardedConcurrentExactlyOneWins 的真并发模式):两个 goroutine 并发对
|
|
// 同一订单投递同一份成功回调,必须恰好一次真正结算 —— 订单 paid 一次(不重复)、
|
|
// webhook outbox 对该单恰好一行(不双投)、一个 goroutine 得 Settled、另一个得 Duplicate。
|
|
//
|
|
// 用真实 store.WebhookStore + webhook.Notifier(非 spyEnqueuer)驱动:本测试要断言的是
|
|
// "落库的 outbox 行数",spyEnqueuer 的内存 slice/map 本身也不是并发安全的,不适合这里。
|
|
// model.OpenTestDB 的 DSN 已带 _txlock=immediate + busy_timeout,并发事务在读阶段即串行化。
|
|
func TestSettleConcurrentSameOrder(t *testing.T) {
|
|
db := openFileGatewaySettleDB(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)
|
|
areg := accounts.New([]config.AccountConfig{
|
|
{AccountID: "fake-a1", Channel: "fake", Region: "global", Enabled: true, Weight: 1},
|
|
})
|
|
picker := accounts.NewRouter(areg, nil, nil)
|
|
ws := store.NewWebhookStore(db)
|
|
notifier := webhook.NewNotifier(ws,
|
|
func(string) (config.BizSystemConfig, bool) { return config.BizSystemConfig{}, false },
|
|
func(string) (bool, error) { return true, nil },
|
|
)
|
|
g := gateway.New(orders, refunds, preg, picker, stubResolver{}, notifier, "global", subs, chargebacks)
|
|
|
|
ctx := context.Background()
|
|
res, err := g.CreateOrder(ctx, gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-conc"})
|
|
if err != nil {
|
|
t.Fatalf("create: %v", err)
|
|
}
|
|
ref := attemptRef(t, orders)
|
|
ev := &provider.PaidEvent{ProviderRef: ref, Status: provider.PaidSucceeded, PaidAmountMinor: 29990000, PaidCurrency: "USDT"}
|
|
|
|
var wg sync.WaitGroup
|
|
results := make([]gateway.SettleResult, 2)
|
|
errs := make([]error, 2)
|
|
for i := 0; i < 2; i++ {
|
|
wg.Add(1)
|
|
go func(i int) {
|
|
defer wg.Done()
|
|
results[i], errs[i] = g.Settle(ctx, ev)
|
|
}(i)
|
|
}
|
|
wg.Wait()
|
|
|
|
for i, gerr := range errs {
|
|
if gerr != nil {
|
|
t.Fatalf("goroutine %d unexpected error: %v", i, gerr)
|
|
}
|
|
}
|
|
processed, duplicate := 0, 0
|
|
for _, r := range results {
|
|
switch r {
|
|
case gateway.SettleProcessed:
|
|
processed++
|
|
case gateway.SettleDuplicate:
|
|
duplicate++
|
|
default:
|
|
t.Fatalf("unexpected result %v in %v", r, results)
|
|
}
|
|
}
|
|
if processed != 1 || duplicate != 1 {
|
|
t.Fatalf("results = %v, want exactly one Processed + one Duplicate", results)
|
|
}
|
|
|
|
o, err := orders.GetOrder(res.OrderNo)
|
|
if err != nil || o.Status != model.OrderPaidV2 {
|
|
t.Fatalf("order = %+v, %v, want paid exactly once", o, err)
|
|
}
|
|
|
|
var cnt int64
|
|
if err := db.Model(&model.WebhookDelivery{}).
|
|
Where("out_trade_no = ? AND event_type = ?", res.OrderNo, "payment.succeeded").
|
|
Count(&cnt).Error; err != nil {
|
|
t.Fatalf("count outbox: %v", err)
|
|
}
|
|
if cnt != 1 {
|
|
t.Fatalf("outbox rows for order = %d, want exactly 1(不双投)", cnt)
|
|
}
|
|
}
|
|
|
|
func TestHandleCallbackAndSync(t *testing.T) {
|
|
g, fp, _, orders := newGateway(t)
|
|
ctx := context.Background()
|
|
g.CreateOrder(ctx, gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-1"})
|
|
ref := attemptRef(t, orders)
|
|
|
|
// 回调路径:fake.VerifyCallback 解析 JSON → Settle
|
|
body, _ := json.Marshal(map[string]any{"provider_ref": ref, "status": "succeeded", "amount_minor": 29990000, "currency": "USDT"})
|
|
got, err := g.HandleCallback(ctx, "fake", provider.CallbackInput{Raw: body})
|
|
if err != nil || got != gateway.SettleProcessed {
|
|
t.Fatalf("HandleCallback = %v, %v", got, err)
|
|
}
|
|
|
|
// 查单兜底:另起一单,预置 query 命中 → SyncPendingAttempts 收敛
|
|
res2, _ := g.CreateOrder(ctx, gateway.CreateOrderInput{SKU: "pro_year", Method: "fake", BizSystem: "pangolin", BizRef: "u-2"})
|
|
atts, _ := orders.ListAttemptsByStatus(model.AttemptPending, 10)
|
|
ref2 := atts[0].ProviderRef
|
|
fp.SetQueryResult(ref2, provider.PaidEvent{ProviderRef: ref2, Status: provider.PaidSucceeded, PaidAmountMinor: 29990000, PaidCurrency: "USDT"})
|
|
n, err := g.SyncPendingAttempts(ctx, 10)
|
|
if err != nil || n < 1 {
|
|
t.Fatalf("SyncPendingAttempts = %d, %v", n, err)
|
|
}
|
|
o2, _ := orders.GetOrder(res2.OrderNo)
|
|
if o2.Status != model.OrderPaidV2 {
|
|
t.Fatalf("查单兜底后 order 应 paid, got %v", o2.Status)
|
|
}
|
|
}
|