From 0ae7769abea5726e8bc81e4cbb2207c178873bff Mon Sep 17 00:00:00 2001 From: wangjia <809946525@qq.com> Date: Sat, 11 Jul 2026 08:50:45 +0800 Subject: [PATCH] =?UTF-8?q?test(pay):=20=E8=A1=A5=20settle=20=E5=B9=B6?= =?UTF-8?q?=E5=8F=91/=E8=B7=A8=E5=8C=85=E5=85=A8=E9=93=BE=20e2e/reconcile?= =?UTF-8?q?=20=E8=A3=85=E9=85=8D/alipay=20=E9=AA=8C=E7=AD=BE=E5=81=A5?= =?UTF-8?q?=E5=A3=AE=E6=80=A7/money=20=E9=9B=B6=E5=B0=8F=E6=95=B0=E4=BD=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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,不是乘除法坑)。 --- internal/gateway/e2e_fullchain_test.go | 151 +++++++++++++++++++ internal/gateway/settle_test.go | 119 +++++++++++++++ internal/money/money_test.go | 33 +++++ internal/provider/alipay/alipay_test.go | 183 ++++++++++++++++++++++++ internal/reconcile/assembly.go | 44 ++++++ internal/reconcile/assembly_test.go | 153 ++++++++++++++++++++ internal/reconcile/runner.go | 10 ++ main.go | 19 +-- 8 files changed, 696 insertions(+), 16 deletions(-) create mode 100644 internal/gateway/e2e_fullchain_test.go create mode 100644 internal/reconcile/assembly.go create mode 100644 internal/reconcile/assembly_test.go diff --git a/internal/gateway/e2e_fullchain_test.go b/internal/gateway/e2e_fullchain_test.go new file mode 100644 index 0000000..3d54195 --- /dev/null +++ b/internal/gateway/e2e_fullchain_test.go @@ -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) + } +} diff --git a/internal/gateway/settle_test.go b/internal/gateway/settle_test.go index 04713c5..525dae8 100644 --- a/internal/gateway/settle_test.go +++ b/internal/gateway/settle_test.go @@ -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() diff --git a/internal/money/money_test.go b/internal/money/money_test.go index 7997554..5ac3680 100644 --- a/internal/money/money_test.go +++ b/internal/money/money_test.go @@ -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 { diff --git a/internal/provider/alipay/alipay_test.go b/internal/provider/alipay/alipay_test.go index 52b8672..1965a9e 100644 --- a/internal/provider/alipay/alipay_test.go +++ b/internal/provider/alipay/alipay_test.go @@ -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() diff --git a/internal/reconcile/assembly.go b/internal/reconcile/assembly.go new file mode 100644 index 0000000..5b08894 --- /dev/null +++ b/internal/reconcile/assembly.go @@ -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 +} diff --git a/internal/reconcile/assembly_test.go b/internal/reconcile/assembly_test.go new file mode 100644 index 0000000..eb3e25b --- /dev/null +++ b/internal/reconcile/assembly_test.go @@ -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) + } +} diff --git a/internal/reconcile/runner.go b/internal/reconcile/runner.go index 31fb76a..08faaf5 100644 --- a/internal/reconcile/runner.go +++ b/internal/reconcile/runner.go @@ -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 { diff --git a/main.go b/main.go index 117a078..9fd1a71 100644 --- a/main.go +++ b/main.go @@ -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) 前的启动