merge: 补测试批B(settle并发/全链e2e/reconcile装配/alipay健壮性/money零小数位)

This commit is contained in:
wangjia
2026-07-11 08:51:35 +08:00
8 changed files with 696 additions and 16 deletions
+151
View File
@@ -0,0 +1,151 @@
package gateway_test
import (
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"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/provider/fake"
"github.com/wangjia/pay/internal/store"
"github.com/wangjia/pay/internal/util"
"github.com/wangjia/pay/internal/webhook"
)
// TestE2EFullChainCallbackToDeliveredWebhook 补齐「下单 → 回调 → settle → webhook
// 实际 HTTP 投递 → 标 delivered」整链(各层此前只独立测过):
// - 起一个 httptest.Server 当业务方回调接收器,记录收到的请求,验证 X-Pay-* 头齐全 +
// 能用同一 secret 验签通过(出站签名互操作,与 internal/webhook/notifier_test.go 同法)。
// - fake 渠道下单 → 模拟渠道成功回调(HandleCallback,同 settle_test.go
// TestHandleCallbackAndSync 的回调体) → Settle 翻转 paid → outbox 入队。
// - 用真实 store.WebhookStore + webhook.Notifier(非 spyEnqueuer)驱动一次
// DeliverPending(即 notifier 的投递 RunOnce)→ 断言 webhook 真被投到 test server、
// server 收到的 payload 签名正确、outbox 行标记 delivered。
func TestE2EFullChainCallbackToDeliveredWebhook(t *testing.T) {
const secret = "e2e-fullchain-secret"
const bizSystem = "pangolin"
var gotBody []byte
var gotHeaders http.Header
var hits int
biz := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
hits++
gotBody, _ = io.ReadAll(r.Body)
gotHeaders = r.Header.Clone()
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("SUCCESS"))
}))
defer biz.Close()
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)
ws := store.NewWebhookStore(db)
// bizConfig:biz_system "pangolin" 的回调地址指向这台 test server(生产由
// config.C.BizByName 按业务方名下发,这里直接指定同等语义)。
bizCfg := func(system string) (config.BizSystemConfig, bool) {
if system == bizSystem {
return config.BizSystemConfig{CallbackURL: biz.URL, Secret: secret}, true
}
return config.BizSystemConfig{}, false
}
// orderPaid 门禁走真实订单查询(不像多数单测那样恒真),验证投递门禁与真实订单状态联动。
orderPaid := func(outTradeNo string) (bool, error) {
o, err := orders.GetOrder(outTradeNo)
if err != nil {
return false, err
}
return o.Status == model.OrderPaidV2, nil
}
notifier := webhook.NewNotifier(ws, bizCfg, orderPaid)
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: bizSystem, BizRef: "u-e2e-full",
})
if err != nil {
t.Fatalf("create order: %v", err)
}
ref := attemptRef(t, orders)
// 渠道方成功回调(同 TestHandleCallbackAndSync 的回调体)→ VerifyCallback → 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, want SettleProcessed", got, err)
}
o, err := orders.GetOrder(res.OrderNo)
if err != nil || o.Status != model.OrderPaidV2 {
t.Fatalf("order = %+v, %v, want paid", o, err)
}
// 入队后未投递:确认 outbox 里躺着一条待发行。
pend, err := ws.ListUndelivered(10)
if err != nil || len(pend) != 1 || pend[0].OutTradeNo != res.OrderNo {
t.Fatalf("undelivered = %+v, %v, want 1 row for %s", pend, err, res.OrderNo)
}
// notifier 投递一轮(RunOnce 语义)→ 真实 HTTP POST 到业务方接收器。
sent, err := notifier.DeliverPending(10)
if err != nil || sent != 1 {
t.Fatalf("DeliverPending = %d, %v, want 1", sent, err)
}
if hits != 1 {
t.Fatalf("业务方应恰收到 1 次 POST, got %d", hits)
}
// 出站签名头齐全 + 能用同一 secret 验签通过(互操作校验,业务方即用同法验签)。
sys := gotHeaders.Get("X-Pay-System")
ev := gotHeaders.Get("X-Pay-Event")
ts := gotHeaders.Get("X-Pay-Timestamp")
nonce := gotHeaders.Get("X-Pay-Nonce")
sign := gotHeaders.Get("X-Pay-Sign")
if sys != bizSystem || ev != "payment.succeeded" || ts == "" || nonce == "" || sign == "" {
t.Fatalf("X-Pay-* 头不全: system=%q event=%q ts=%q nonce=%q sign=%q", sys, ev, ts, nonce, sign)
}
if !util.HMACVerify(secret, sign, sys, ts, nonce, string(gotBody)) {
t.Fatalf("业务方侧验签失败(出站签名与 HMACVerify 不互操作)")
}
var payload map[string]any
if err := json.Unmarshal(gotBody, &payload); err != nil {
t.Fatalf("payload 非法 JSON: %v, body=%s", err, gotBody)
}
if payload["event_type"] != "payment.succeeded" || payload["out_trade_no"] != res.OrderNo ||
payload["product_biz_code"] != "pro_year" {
t.Fatalf("payload = %+v", payload)
}
// outbox 行已标 delivered:再投一次不应重发。
pendAfter, err := ws.ListUndelivered(10)
if err != nil || len(pendAfter) != 0 {
t.Fatalf("投递后 undelivered = %+v, %v, want 0(已标 delivered)", pendAfter, err)
}
if sent2, _ := notifier.DeliverPending(10); sent2 != 0 {
t.Fatalf("已投递不应重发, got %d", sent2)
}
if hits != 1 {
t.Fatalf("重投检查不应产生新 POST, hits=%d", hits)
}
}
+119
View File
@@ -3,8 +3,14 @@ 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"
@@ -12,6 +18,7 @@ import (
"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 {
@@ -169,6 +176,118 @@ func TestSettleTransientReadErrorIsFailed(t *testing.T) {
}
}
// 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()
+33
View File
@@ -67,6 +67,39 @@ func TestParseOverflow(t *testing.T) {
}
}
// TestZeroDecimalCurrencyJPYNotSupported 覆盖补测项 ⑩a:先读 internal/money/money.go
// 确认现状——exponents 只注册了 CNY/USD/USDT(2/2/6 位小数),不含任何 0 位小数币种
// (JPY/KRW 等)。money.Format 内部虽然写了 `if exp == 0 { return strconv.FormatInt(...) }`
// 分支(不做除法/拼小数点),说明"按币种可变小数位"这个架构本身具备支持零小数位
// 的能力,Parse 侧 pow10Int64(0)==1 也天然不做 *100/÷100 那种隐式换算;但因为
// exponents 表里从未注册过任何 exp=0 的币种,这条分支在生产/测试代码里从未被真正
// 触发过——**JPY 本身当前不支持**,接入会在 Exponent()/Parse()/Format() 三处直接拿到
// ErrUnknownCurrency(是"未知币种"硬拒,不是更隐蔽的"用错乘除法"那种坑)。
//
// 待改进记录(供后续接入 JPY/KRW 时参考,未经验证的推测,不代表已确认可行):
// 读代码看,大概率只需在 money.go 的 exponents map 里补一行 "JPY": 0 / "KRW": 0,
// Format 的 exp==0 分支、Parse 的 pow10Int64(0)=1 应该不用改——但这只是读源码后的
// 推断,本测试的黑盒范围(money_test 外部包,不碰 money.go)验证不了"补一行就真的
// 够了",只钉住"今天 JPY/KRW 确实不支持"这一现状,防止有人以为已经支持了。
func TestZeroDecimalCurrencyJPYNotSupported(t *testing.T) {
if _, ok := money.Exponent("JPY"); ok {
t.Fatal("JPY 已被注册进 exponents——本测试钉的是「未注册」旧行为,请更新/删除,并按其真实小数位数补 Parse/Format round-trip 用例")
}
if _, err := money.Parse("100", "JPY"); !errors.Is(err, money.ErrUnknownCurrency) {
t.Fatalf("JPY 未注册,Parse 应报 ErrUnknownCurrency, got %v", err)
}
if _, err := money.Format(100, "JPY"); !errors.Is(err, money.ErrUnknownCurrency) {
t.Fatalf("JPY 未注册,Format 应报 ErrUnknownCurrency, got %v", err)
}
// KRW(韩元)同为真实零小数位币种,同样未注册——顺带钉住,避免误以为只有 JPY 漏了。
if _, ok := money.Exponent("KRW"); ok {
t.Fatal("KRW 已被注册进 exponents——本测试钉的是「未注册」旧行为,请更新/删除")
}
if _, err := money.Parse("100", "KRW"); !errors.Is(err, money.ErrUnknownCurrency) {
t.Fatalf("KRW 未注册,Parse 应报 ErrUnknownCurrency, got %v", err)
}
}
func TestParseFormatNegative(t *testing.T) {
minor, err := money.Parse("-12.34", "CNY")
if err != nil {
+183
View File
@@ -178,6 +178,189 @@ func TestVerifyCallbackRSA_ClosedIsFailed(t *testing.T) {
}
}
// ---------------------------------------------------------------------------
// 健壮性测试(补测项 ⑨):以上验签测试都是自生成密钥对自签自验(测的是我方 adapter
// 逻辑本身,不是支付宝真实签名细节);本节仍是自签,但专挑真实回调报文常见的
// "互操作易碎点"——未知/多余字段、非字母序的字段排列、中文/unicode 值、可选字段
// 留空。⚠️ 真实支付宝黄金向量(官方真实回调报文的逐字节样本)需要抓真实线上/沙箱
// 回调抓包,本文件不具备,以下均是自签构造。
// ---------------------------------------------------------------------------
// TestVerifyCallbackRSA_UnknownExtraFieldsIgnored 真实支付宝回调常带一大堆我方代码
// 不读取的字段(app_id/seller_email/buyer_logon_id/notify_id/charset/version/
// receipt_amount/point_amount...)。signRSA2 与 SDK 内部验签都按官方规则排除
// sign/sign_type 及空值字段参与签名重算,多余字段不应打断验签,也不应污染归一化
// 结果(ProviderRef/Status 只认 out_trade_no/trade_status/total_amount)。
func TestVerifyCallbackRSA_UnknownExtraFieldsIgnored(t *testing.T) {
appPriv, aliPriv, aliPub := genKeys(t)
p := ali.New(buildClient(t, appPriv, aliPub))
form := url.Values{}
form.Set("out_trade_no", "PAY-EXTRA")
form.Set("trade_no", "2021EXTRA")
form.Set("trade_status", "TRADE_SUCCESS")
form.Set("total_amount", "1.00")
form.Set("sign_type", "RSA2")
// 支付宝真实通知常见的额外字段,adapter 均不读取,只应被忽略。
form.Set("app_id", "2021000000000000")
form.Set("auth_app_id", "2021000000000000")
form.Set("notify_id", "notify-xyz")
form.Set("notify_type", "trade_status_sync")
form.Set("notify_time", "2026-07-10 15:04:05")
form.Set("charset", "utf-8")
form.Set("version", "1.0")
form.Set("seller_id", "2088seller")
form.Set("seller_email", "seller@example.com")
form.Set("buyer_id", "2088buyer")
form.Set("buyer_logon_id", "buy***@example.com")
form.Set("receipt_amount", "1.00")
form.Set("point_amount", "0.00")
form.Set("sign", signRSA2(t, aliPriv, form))
ev, err := p.VerifyCallback(context.Background(), provider.CallbackInput{Raw: []byte(form.Encode())})
if err != nil {
t.Fatalf("verify: %v", err)
}
if ev.ProviderRef != "PAY-EXTRA" || ev.Status != provider.PaidSucceeded || ev.PaidAmountMinor != 100 {
t.Fatalf("event = %+v", ev)
}
}
// TestVerifyCallbackRSA_FieldOrderIndependent 真实网关/CDN/负载均衡转发不保证 query
// 字段在线路上的排列顺序;VerifyCallback 先 url.ParseQuery 解成 url.Values(map),
// 理论上字段顺序不该影响结果——这里故意手拼一份与字母序相反的 wire bytes(而非用
// form.Encode() 那种天然按 key 字母序输出的路径)来钉住这一互操作预期。
func TestVerifyCallbackRSA_FieldOrderIndependent(t *testing.T) {
appPriv, aliPriv, aliPub := genKeys(t)
p := ali.New(buildClient(t, appPriv, aliPub))
form := url.Values{}
form.Set("out_trade_no", "PAY-ORDER")
form.Set("trade_no", "2021ORDER")
form.Set("trade_status", "TRADE_SUCCESS")
form.Set("total_amount", "1.00")
form.Set("gmt_payment", "2026-07-10 15:04:05")
form.Set("sign_type", "RSA2")
sign := signRSA2(t, aliPriv, form) // 按官方排序算出的签名,与 wire 上字段实际排列顺序无关
raw := "sign_type=" + url.QueryEscape("RSA2") +
"&trade_status=" + url.QueryEscape("TRADE_SUCCESS") +
"&total_amount=" + url.QueryEscape("1.00") +
"&gmt_payment=" + url.QueryEscape("2026-07-10 15:04:05") +
"&out_trade_no=" + url.QueryEscape("PAY-ORDER") +
"&trade_no=" + url.QueryEscape("2021ORDER") +
"&sign=" + url.QueryEscape(sign)
ev, err := p.VerifyCallback(context.Background(), provider.CallbackInput{Raw: []byte(raw)})
if err != nil {
t.Fatalf("verify(乱序字段): %v", err)
}
if ev.ProviderRef != "PAY-ORDER" || ev.Status != provider.PaidSucceeded {
t.Fatalf("event = %+v", ev)
}
if ev.PaidAt == nil {
t.Fatal("乱序字段不应影响 gmt_payment 解析,PaidAt 不应为 nil")
}
}
// TestVerifyCallbackRSA_UnicodeSubject subject/body 携带中文 + emoji(真实商品名称/
// 备注常见),须能正确 percent-encode/decode 并参与签名而不破坏验签或后续解析。
func TestVerifyCallbackRSA_UnicodeSubject(t *testing.T) {
appPriv, aliPriv, aliPub := genKeys(t)
p := ali.New(buildClient(t, appPriv, aliPub))
form := url.Values{}
form.Set("out_trade_no", "PAY-CN")
form.Set("trade_no", "2021CN")
form.Set("trade_status", "TRADE_SUCCESS")
form.Set("total_amount", "1.00")
form.Set("subject", "Pro 年付 🎉 商品名称含中文与emoji")
form.Set("body", "订单备注:测试中文正文——含标点、全角符号「」")
form.Set("sign_type", "RSA2")
form.Set("sign", signRSA2(t, aliPriv, form))
ev, err := p.VerifyCallback(context.Background(), provider.CallbackInput{Raw: []byte(form.Encode())})
if err != nil {
t.Fatalf("verify: %v", err)
}
if ev.ProviderRef != "PAY-CN" || ev.Status != provider.PaidSucceeded || ev.PaidAmountMinor != 100 {
t.Fatalf("event = %+v", ev)
}
}
// TestVerifyCallbackRSA_EmptyValueFieldsAreIncludedInSignature 补测中意外挖出的真实
// SDK 行为,与"官方文档说空值参数不参与签名"的字面印象不一致,记录在案(避免下次
// 被当成 bug 重新踩一遍):github.com/smartwalle/alipay/v3 注入给签名器的是它自己
// vendor 的 Encoder(v3@v3.2.29/encode.go),不是 nsign 包默认的 DefaultEncoder——
// 前者对"空值字段"**不做排除**,只要 key 不在 ignore 名单(sign/sign_type/
// alipay_cert_sn)里,即便 value 是空字符串也会以 "key=" 形式进入签名源参与 hash;
// nsign.DefaultEncoder 则会跳过空值。
//
// 实测过程:先按"官方文档描述的排除空值"假设写了个用例(签名时排除空值、wire 上
// 带 3 个空值可选字段),结果 SDK 侧验签失败(crypto/rsa: verification error);手工
// 复算签名源字符串确认两侧字符串其实一致后,才定位到问题出在 vendor Encoder 与
// nsign.DefaultEncoder 的语义差异,而非我方 adapter 或字段解析逻辑的 bug。
//
// 本用例改为验证 SDK 的**真实行为**:空值字段必须参与签名源才能验签通过(signRSA2
// 复刻的是"排除空值"这一支付宝官方文档口径,与此处刻意用 signIncludingEmpty 区分)。
// 真实支付宝服务器是否会在回调里真的带上空值字段未经抓包验证,但只要它带了,
// 当前这版 SDK 要求空值也进签名源——这条钉住的是 SDK 版本的真实契约,供后续升级
// SDK 或换实现时对照。
func TestVerifyCallbackRSA_EmptyValueFieldsAreIncludedInSignature(t *testing.T) {
appPriv, aliPriv, aliPub := genKeys(t)
p := ali.New(buildClient(t, appPriv, aliPub))
form := url.Values{}
form.Set("out_trade_no", "PAY-EMPTY")
form.Set("trade_no", "2021EMPTY")
form.Set("trade_status", "TRADE_SUCCESS")
form.Set("total_amount", "1.00")
form.Set("sign_type", "RSA2")
form.Set("refund_status", "")
form.Set("invoice_amount", "")
form.Set("passback_params", "")
form.Set("sign", signIncludingEmpty(t, aliPriv, form))
ev, err := p.VerifyCallback(context.Background(), provider.CallbackInput{Raw: []byte(form.Encode())})
if err != nil {
t.Fatalf("verify(含空值可选字段): %v", err)
}
if ev.ProviderRef != "PAY-EMPTY" || ev.Status != provider.PaidSucceeded {
t.Fatalf("event = %+v", ev)
}
}
// signIncludingEmpty 复刻 alipay 包 vendor Encoder(见上面测试注释,非 nsign 默认
// DefaultEncoder)的真实签名规则:排序**全部**非 sign/sign_type 字段(空值也算入),
// k=v&拼接,RSA-SHA256,base64。与本文件 signRSA2(严格照官方文档排除空值)刻意
// 区分,专供本测试还原 SDK 真实行为,不共用逻辑以免掩盖两者的语义差异。
func signIncludingEmpty(t *testing.T, aliPrivB64 string, form url.Values) string {
t.Helper()
der, _ := base64.StdEncoding.DecodeString(aliPrivB64)
priv, err := x509.ParsePKCS1PrivateKey(der)
if err != nil {
t.Fatalf("parse ali priv: %v", err)
}
keys := make([]string, 0, len(form))
for k := range form {
if k == "sign" || k == "sign_type" {
continue
}
keys = append(keys, k)
}
sort.Strings(keys)
var parts []string
for _, k := range keys {
parts = append(parts, k+"="+form.Get(k))
}
h := sha256.Sum256([]byte(strings.Join(parts, "&")))
sig, err := rsa.SignPKCS1v15(rand.Reader, priv, crypto.SHA256, h[:])
if err != nil {
t.Fatalf("sign: %v", err)
}
return base64.StdEncoding.EncodeToString(sig)
}
// signRSA2 复刻支付宝签名:排序非空参数(排除 sign/sign_type),k=v&拼接,RSA-SHA256,base64。
func signRSA2(t *testing.T, aliPrivB64 string, form url.Values) string {
t.Helper()
+44
View File
@@ -0,0 +1,44 @@
package reconcile
import (
"time"
"github.com/wangjia/pay/config"
"github.com/wangjia/pay/internal/accounts"
"github.com/wangjia/pay/internal/gateway"
"github.com/wangjia/pay/internal/provider"
"github.com/wangjia/pay/internal/store"
)
// Assemble builds the full P6 background-reconcile Runner: order-expire /
// usage-refresh / refund-apply-sweep / refund-stuck-alert (local, DB-only,
// AddLocal — run synchronously once at startup by main.go's RunOnceLocal) +
// sync-pending / paid-spotcheck (network, Add — first tick handles it) +
// crypto warm(直跑)+ crypto-orphan-scan(AddCryptoJobs,network,只在 crypto 渠道
// 已注册时挂载).
//
// 抽出这个函数纯粹是为了让"预期任务集合是否都注册上"能被直接单测(main.go 原先把
// 这段写死在 main() 里,main() 本身因为要拉真实 DB/HTTP server 不适合单测——见
// assembly_test.go)。main.go 对应 `if config.C.Reconcile.Enabled` 分支现在只是薄薄
// 一层:调这里 + RunOnceLocal/Start/日志,行为与抽取前逐行一致,未改任何调度语义。
func Assemble(orderStore *store.OrderStore, refundStore *store.RefundStore, usage *UsageSource,
gw *gateway.Gateway, pReg *provider.Registry, acctReg *accounts.Registry,
orphanStore *store.OrphanStore, rc config.ReconcileConfig) *Runner {
runner := NewRunner()
// 纯本地(DB-only)任务用 AddLocal:会被启动预热 RunOnceLocal 同步跑一遍。
runner.AddLocal("order-expire", time.Duration(rc.ExpireEverySec)*time.Second,
OrderExpirerTask(orderStore, time.Duration(rc.OrderTTLMin)*time.Minute, time.Now))
runner.AddLocal("usage-refresh", time.Duration(rc.UsageEverySec)*time.Second,
RefreshUsageTask(usage))
runner.AddLocal("refund-apply-sweep", time.Duration(rc.RefundApplyEverySec)*time.Second,
RefundApplyTask(orderStore, refundStore, time.Duration(rc.RefundApplyLookbackMin)*time.Minute, time.Now, 200))
runner.AddLocal("refund-stuck-alert", time.Duration(rc.RefundApplyEverySec)*time.Second,
RefundStuckAlertTask(refundStore, time.Duration(rc.RefundStuckWarnMin)*time.Minute, time.Now))
// 网络型(出网 HTTP,10-15s 超时)任务用 Add:启动预热不跑,交各自 ticker 首跳。
runner.Add("sync-pending", time.Duration(rc.SyncEverySec)*time.Second,
SyncPendingTask(gw, 100))
runner.Add("paid-spotcheck", time.Duration(rc.SpotCheckEverySec)*time.Second,
PaidSpotCheckTask(orderStore, pReg, time.Duration(rc.SpotCheckWindowMin)*time.Minute, time.Now))
AddCryptoJobs(runner, pReg, acctReg, orderStore, orphanStore, rc) // crypto 预留冷启动 Warm(直跑,不受影响)+ 孤儿扫描(网络型,交 ticker 首跳)
return runner
}
+153
View File
@@ -0,0 +1,153 @@
package reconcile_test
import (
"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/provider/crypto"
"github.com/wangjia/pay/internal/provider/fake"
"github.com/wangjia/pay/internal/reconcile"
"github.com/wangjia/pay/internal/store"
)
// TestAssembleRegistersAllExpectedTasks 钉住 main.go 里那一长串 reconcile 任务(P6
// order-expire/usage-refresh/refund-apply-sweep/refund-stuck-alert/sync-pending/
// paid-spotcheck + P8 Task6 crypto-orphan-scan)确实全部被注册 —— main() 原先把这段
// 装配写死在函数体内,没有任何测试保证"漏挂一个任务"能被发现;抽成
// reconcile.Assemble(见 internal/reconcile/assembly.go)后可以直接断言注册出的任务名
// 集合,main.go 现在只是薄薄一层调用(逐行照抄原实现,未改调度语义)。
//
// pReg 同时注册 fake + crypto 两个渠道:crypto-orphan-scan 只在 crypto 渠道已注册时
// 才由 AddCryptoJobs 挂载(见其函数注释),不带 crypto 就测不到这第 7 个任务。
func TestAssembleRegistersAllExpectedTasks(t *testing.T) {
db := model.OpenTestDB(t)
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
orphans := store.NewOrphanStore(db)
preg := provider.NewRegistry()
fp := fake.New()
preg.Register(fp)
acctReg := accounts.New([]config.AccountConfig{
{AccountID: "fake-a1", Channel: "fake", Region: "global", Enabled: true},
{AccountID: "crypto-a1", Channel: "crypto", Region: "global", Enabled: true, CredentialEnvPrefix: "assembletest"},
})
preg.Register(crypto.New(acctReg))
picker := accounts.NewRouter(acctReg, nil, nil)
usage := reconcile.NewUsageSource(orders, nil)
gw := gateway.New(orders, refunds, preg, picker, stubResolver{}, nopEnq{}, "global", subs, chargebacks)
rc := config.ReconcileConfig{
Enabled: true, OrderTTLMin: 30, ExpireEverySec: 60, SyncEverySec: 60, UsageEverySec: 60,
SpotCheckEverySec: 60, SpotCheckWindowMin: 60, OrphanEverySec: 60, OrphanWindowMin: 60,
RefundApplyEverySec: 60, RefundStuckWarnMin: 60, RefundApplyLookbackMin: 60,
}
runner := reconcile.Assemble(orders, refunds, usage, gw, preg, acctReg, orphans, rc)
want := []string{
"order-expire", "usage-refresh", "refund-apply-sweep", "refund-stuck-alert",
"sync-pending", "paid-spotcheck", "crypto-orphan-scan",
}
got := runner.TaskNames()
if len(got) != len(want) {
t.Fatalf("task names = %v (%d), want %d: %v", got, len(got), len(want), want)
}
seen := make(map[string]bool, len(got))
for _, n := range got {
seen[n] = true
}
for _, w := range want {
if !seen[w] {
t.Fatalf("缺任务 %q,got %v", w, got)
}
}
}
// TestAssembleSkipsCryptoOrphanScanWithoutCryptoChannel crypto 渠道未启用时(providerbuild
// 按账户配置决定是否注册,见 internal/providerbuild),AddCryptoJobs 安全跳过——不应
// 注册 crypto-orphan-scan,其余 6 个本地/网络任务照常注册,总数少一个。
func TestAssembleSkipsCryptoOrphanScanWithoutCryptoChannel(t *testing.T) {
db := model.OpenTestDB(t)
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
orphans := store.NewOrphanStore(db)
preg := provider.NewRegistry() // 只 fake,没有 crypto
fp := fake.New()
preg.Register(fp)
acctReg := accounts.New([]config.AccountConfig{
{AccountID: "fake-a1", Channel: "fake", Region: "global", Enabled: true},
})
picker := accounts.NewRouter(acctReg, nil, nil)
usage := reconcile.NewUsageSource(orders, nil)
gw := gateway.New(orders, refunds, preg, picker, stubResolver{}, nopEnq{}, "global", subs, chargebacks)
rc := config.ReconcileConfig{
Enabled: true, OrderTTLMin: 30, ExpireEverySec: 60, SyncEverySec: 60, UsageEverySec: 60,
SpotCheckEverySec: 60, SpotCheckWindowMin: 60, OrphanEverySec: 60, OrphanWindowMin: 60,
RefundApplyEverySec: 60, RefundStuckWarnMin: 60, RefundApplyLookbackMin: 60,
}
runner := reconcile.Assemble(orders, refunds, usage, gw, preg, acctReg, orphans, rc)
got := runner.TaskNames()
if len(got) != 6 {
t.Fatalf("无 crypto 渠道时应只注册 6 个任务(缺 crypto-orphan-scan), got %v", got)
}
for _, n := range got {
if n == "crypto-orphan-scan" {
t.Fatalf("未启用 crypto 渠道时不应注册 crypto-orphan-scan, got %v", got)
}
}
}
// TestAssembleZeroIntervalSkipsThatTaskOnly 覆盖 Runner.add 的 interval<=0 守卫在装配层
// 的联动:某一项周期配成 0(operator 配置疏漏,如 usage_every_sec: 0)只会让那一个任务
// 不注册,其余任务不受影响(呼应 runner_test.go 的
// TestRunnerZeroIntervalSkipsRegistrationAndStartDoesNotPanic,但那里是直接摆弄 Runner,
// 这里钉住 Assemble 真按 config 字段逐个转发 interval,没有哪个任务共用错了字段)。
func TestAssembleZeroIntervalSkipsThatTaskOnly(t *testing.T) {
db := model.OpenTestDB(t)
orders := store.NewOrderStore(db)
refunds := store.NewRefundStore(db)
subs := store.NewSubscriptionStore(db)
chargebacks := store.NewChargebackStore(db)
orphans := store.NewOrphanStore(db)
preg := provider.NewRegistry()
fp := fake.New()
preg.Register(fp)
acctReg := accounts.New([]config.AccountConfig{{AccountID: "fake-a1", Channel: "fake", Region: "global", Enabled: true}})
picker := accounts.NewRouter(acctReg, nil, nil)
usage := reconcile.NewUsageSource(orders, nil)
gw := gateway.New(orders, refunds, preg, picker, stubResolver{}, nopEnq{}, "global", subs, chargebacks)
rc := config.ReconcileConfig{
Enabled: true, OrderTTLMin: 30, ExpireEverySec: 60, SyncEverySec: 60,
UsageEverySec: 0, // 疏漏:只这一项配成 0
SpotCheckEverySec: 60, SpotCheckWindowMin: 60, OrphanEverySec: 60, OrphanWindowMin: 60,
RefundApplyEverySec: 60, RefundStuckWarnMin: 60, RefundApplyLookbackMin: 60,
}
runner := reconcile.Assemble(orders, refunds, usage, gw, preg, acctReg, orphans, rc)
got := runner.TaskNames()
for _, n := range got {
if n == "usage-refresh" {
t.Fatalf("usage_every_sec=0 时 usage-refresh 不应注册, got %v", got)
}
}
// 其余 5 个(order-expire/refund-apply-sweep/refund-stuck-alert/sync-pending/paid-spotcheck)照常。
if len(got) != 5 {
t.Fatalf("只 usage-refresh 应被跳过,其余应全注册, got %v", got)
}
}
+10
View File
@@ -86,6 +86,16 @@ func (r *Runner) RunOnceLocal(ctx context.Context) {
}
}
// TaskNames 返回已注册任务名(注册序)。供装配期测试断言"预期任务集合是否都注册上"
// (见 assembly.go::Assemble 与 assembly_test.go),不用于运行期逻辑。
func (r *Runner) TaskNames() []string {
names := make([]string, len(r.tasks))
for i, t := range r.tasks {
names[i] = t.Name
}
return names
}
// Start 每任务一 goroutine + 独立 ticker 常驻;ctx 取消即退出。每 tick 崩溃安全。
func (r *Runner) Start(ctx context.Context) {
for _, t := range r.tasks {
+3 -16
View File
@@ -108,23 +108,10 @@ func main() {
// crypto 预留冷启动 Warm + 孤儿扫描(Task 6,AddCryptoJobs)。
if config.C.Reconcile.Enabled {
rc := config.C.Reconcile
runner := reconcile.NewRunner()
// 纯本地(DB-only)任务用 AddLocal:会被启动预热 RunOnceLocal 同步跑一遍。
runner.AddLocal("order-expire", time.Duration(rc.ExpireEverySec)*time.Second,
reconcile.OrderExpirerTask(orderStore, time.Duration(rc.OrderTTLMin)*time.Minute, time.Now))
runner.AddLocal("usage-refresh", time.Duration(rc.UsageEverySec)*time.Second,
reconcile.RefreshUsageTask(usage))
runner.AddLocal("refund-apply-sweep", time.Duration(rc.RefundApplyEverySec)*time.Second,
reconcile.RefundApplyTask(orderStore, refundStore, time.Duration(rc.RefundApplyLookbackMin)*time.Minute, time.Now, 200))
runner.AddLocal("refund-stuck-alert", time.Duration(rc.RefundApplyEverySec)*time.Second,
reconcile.RefundStuckAlertTask(refundStore, time.Duration(rc.RefundStuckWarnMin)*time.Minute, time.Now))
// 网络型(出网 HTTP,10-15s 超时)任务用 Add:启动预热不跑,交各自 ticker 首跳。
runner.Add("sync-pending", time.Duration(rc.SyncEverySec)*time.Second,
reconcile.SyncPendingTask(gw, 100))
runner.Add("paid-spotcheck", time.Duration(rc.SpotCheckEverySec)*time.Second,
reconcile.PaidSpotCheckTask(orderStore, pReg, time.Duration(rc.SpotCheckWindowMin)*time.Minute, time.Now))
orphanStore := store.NewOrphanStore(db)
reconcile.AddCryptoJobs(runner, pReg, acctReg, orderStore, orphanStore, rc) // crypto 预留冷启动 Warm(直跑,不受影响)+ 孤儿扫描(网络型,交 ticker 首跳)
// 装配逻辑抽到 reconcile.Assemble(单测直接断言任务集合,见
// internal/reconcile/assembly_test.go);这里只是薄薄一层调用 + 预热/常驻启动。
runner := reconcile.Assemble(orderStore, refundStore, usage, gw, pReg, acctReg, orphanStore, rc)
ctx := context.Background()
runner.RunOnceLocal(ctx) // 启动预热:只跑本地任务(usage 快照/过期清理/退款自愈立即生效);网络型任务(sync-pending/paid-spotcheck/crypto-orphan-scan)交各自 ticker 首跳,避免上游慢拖住 r.Run(addr) 前的启动