diff --git a/server/cmd/server/main.go b/server/cmd/server/main.go index fd9e54b..fda8e82 100644 --- a/server/cmd/server/main.go +++ b/server/cmd/server/main.go @@ -379,6 +379,9 @@ func mountV1(r chi.Router, sqlDB *sql.DB, rdb *redis.Client, nodeSvc *nodes.Serv paySystem, paySecret, 5*time.Minute, 15*time.Minute) payWebhook.SetRewarder(rewardSvc) payWebhook.SetNoticer(noticesStore) + // reconcile-on-read:查单发现 pay 网关已 paid 但本地未开通时,复用 webhook 的 + // 幂等 settle 就地开通,不干等 webhook 送达(webhook 仍作兜底)。 + payHandler.SetSettler(payWebhook) } else { log.Printf("PAY_BASE_URL 未配置 — /v1/pay 支付端点不挂载") } diff --git a/server/internal/pay/handler.go b/server/internal/pay/handler.go index 93d6fb6..cd919dd 100644 --- a/server/internal/pay/handler.go +++ b/server/internal/pay/handler.go @@ -1,6 +1,7 @@ package pay import ( + "context" "database/sql" "encoding/json" "errors" @@ -15,18 +16,34 @@ import ( "github.com/wangjia/pangolin/server/internal/auth" ) +// orderSettler 抽象 webhook 的幂等开通入口(settle),供查单对账 reconcile-on-read +// 复用。*WebhookHandler 已满足。放接口而非直接持有 *WebhookHandler,便于测试注入。 +type orderSettler interface { + settle(ctx context.Context, ev *webhookEvent) error +} + +// payOrderPaid 是 pay 网关订单「已支付」状态词(= pay 侧 model.OrderStatusV2 "paid", +// 经 client.GetOrder 原样透传)。查单见此即视为款已到、可开通。改这里前先对齐 pay 契约。 +const payOrderPaid = "paid" + // Handler 是面向 App 的下单代理(JWT 保护;user→biz_ref 映射在 server 侧, // 客户端只传 sku+method+端型 metadata,永远不传金额)。 type Handler struct { - client *Client - store *Store - db *sql.DB + client *Client + store *Store + db *sql.DB + settler orderSettler // 可选:reconcile-on-read 开通入口(见 SetSettler),nil 则退回等 webhook } func NewHandler(client *Client, store *Store, db *sql.DB) *Handler { return &Handler{client: client, store: store, db: db} } +// SetSettler 挂载对账开通入口(reconcile-on-read):装配后 GetOrder 在发现 pay 网关 +// 已 paid 但本地台账未开通时,就地幂等 settle 立即开通,不干等 webhook 送达 +// (webhook 仍作兜底)。为 nil 时退回旧行为(只认本地 row.Status)。 +func (h *Handler) SetSettler(s orderSettler) { h.settler = s } + // allowedMetadataKeys 与 pay gateway.go 白名单一致(is_mobile/render)。 var allowedMetadataKeys = map[string]bool{"is_mobile": true, "render": true} @@ -252,6 +269,30 @@ func (h *Handler) GetOrder(w http.ResponseWriter, r *http.Request) { writePayErr(w, err) return } + // reconcile-on-read:pay 网关已 paid 但本地台账还没开通(webhook 未到/漏投)→ 就地 + // 幂等 settle 立即开通,把「等 webhook 送达」那十几秒砍成「付完下一拍查单即开通」。 + // settle 已 paid 会短路(不会重复开通/发奖);失败仅记日志、退回旧展示,由下一拍轮询 + // 或 webhook 兜底。channel 查单不返回 → 回退 row.Method(仅影响渠道显示,不影响开通)。 + if h.settler != nil && st.Status == payOrderPaid && row.Status != "paid" { + paidAt := "" + if st.PaidAt != nil { + paidAt = st.PaidAt.UTC().Format(time.RFC3339) + } + if serr := h.settler.settle(ctx, &webhookEvent{ + EventType: "payment.succeeded", + OutTradeNo: orderNo, + BizRef: row.BizRef, + ProductBizCode: row.SKU, + AmountMinor: st.AmountMinor, + Currency: st.Currency, + Channel: row.Method, + PaidAt: paidAt, + }); serr != nil { + slog.Warn("pay: reconcile-on-read 开通失败(webhook 仍兜底)", "order_no", orderNo, "err", serr) + } else if r2, e2 := h.store.GetForUser(ctx, uid, orderNo); e2 == nil { + row = r2 // 反映刚开通:本拍即回 activated=true,免再等一拍 + } + } resp := orderStatusResponse{orderView: orderViewFromRow(row), PayStatus: st.Status, Activated: row.Status == "paid"} if row.SubID.Valid { if exp, err := h.store.SubscriptionExpiry(ctx, row.SubID.Int64); err == nil { diff --git a/server/internal/pay/handler_reconcile_sqlite_test.go b/server/internal/pay/handler_reconcile_sqlite_test.go new file mode 100644 index 0000000..dd72137 --- /dev/null +++ b/server/internal/pay/handler_reconcile_sqlite_test.go @@ -0,0 +1,107 @@ +package pay + +import ( + "context" + "database/sql" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/go-chi/chi/v5" + + "github.com/wangjia/pangolin/server/internal/codes" +) + +// reconcileRig:假 pay(查单可配)+ sqlite 台账 + 真 WebhookHandler 作 settler。 +// 验证 GetOrder 的 reconcile-on-read:pay 网关已 paid 时就地幂等开通,不干等 webhook。 +func reconcileRig(t *testing.T, payFn http.HandlerFunc) (*chi.Mux, *Store, *sql.DB) { + t.Helper() + db := openMigratedSQLite(t) + seedUser(t, db, 1, "uuid-1") + srv := fakePay(t, payFn) + st := NewStore(db) + h := NewHandler(NewClient(srv.URL, "pangolin", testSecret), st, db) + // 真 granter + webhook,settle 会真的开通订阅(与线上同一段逻辑)。 + codesSvc := codes.NewService(codes.NewStore(db), nil, 5, time.Hour) + wh := NewWebhookHandler(st, codesSvc, db, nil, "pangolin", testSecret, 5*time.Minute, 15*time.Minute) + h.SetSettler(wh) + r := chi.NewRouter() + r.Get("/v1/pay/orders/{orderNo}", h.GetOrder) + return r, st, db +} + +func getOrder(t *testing.T, router *chi.Mux, orderNo string, uid int64) (int, bool) { + t.Helper() + w := httptest.NewRecorder() + router.ServeHTTP(w, authed(httptest.NewRequest(http.MethodGet, "/v1/pay/orders/"+orderNo, nil), uid)) + var resp struct { + Activated bool `json:"activated"` + } + _ = json.Unmarshal(w.Body.Bytes(), &resp) + return w.Code, resp.Activated +} + +// 网关已 paid、本地还 created → 查单即就地开通(activated=true、台账翻 paid、订阅真授予); +// 且重复轮询幂等:再查一次不二次开通(仍 1 条订阅)。 +func TestGetOrder_ReconcileOnRead_ActivatesWhenGatewayPaid(t *testing.T) { + router, st, db := reconcileRig(t, func(w http.ResponseWriter, _ *http.Request) { + _, _ = w.Write([]byte(`{"data":{"order_no":"pay001","status":"paid", + "subject":"Pro","amount_minor":2999,"currency":"CNY"}}`)) + }) + ctx := context.Background() + if err := st.Insert(ctx, 1, "uuid-1", "pro_month", "pay001", "alipay", 2999, "CNY"); err != nil { + t.Fatal(err) + } + + code, activated := getOrder(t, router, "pay001", 1) + if code != http.StatusOK || !activated { + t.Fatalf("首查应就地开通:code=%d activated=%v", code, activated) + } + + row, err := st.GetForUser(ctx, 1, "pay001") + if err != nil || row.Status != "paid" { + t.Fatalf("本地台账应已翻 paid: row=%+v err=%v", row, err) + } + var n int + if err := db.QueryRow(`SELECT COUNT(*) FROM subscriptions WHERE user_id=1 AND source='pay'`).Scan(&n); err != nil { + t.Fatal(err) + } + if n != 1 { + t.Fatalf("应授予 1 条 pay 订阅,得 %d", n) + } + + // 幂等:再查一次(轮询会持续查),不得二次开通。 + if _, activated2 := getOrder(t, router, "pay001", 1); !activated2 { + t.Fatal("二次查询仍应 activated") + } + _ = db.QueryRow(`SELECT COUNT(*) FROM subscriptions WHERE user_id=1 AND source='pay'`).Scan(&n) + if n != 1 { + t.Fatalf("重复轮询不得二次开通,订阅数 = %d, want 1", n) + } +} + +// 网关仍 pending → 不开通(activated=false、无订阅);证明只有 paid 才触发对账。 +func TestGetOrder_ReconcileOnRead_SkipsWhenGatewayPending(t *testing.T) { + router, st, db := reconcileRig(t, func(w http.ResponseWriter, _ *http.Request) { + _, _ = w.Write([]byte(`{"data":{"order_no":"pay001","status":"pending","currency":"CNY"}}`)) + }) + ctx := context.Background() + if err := st.Insert(ctx, 1, "uuid-1", "pro_month", "pay001", "alipay", 2999, "CNY"); err != nil { + t.Fatal(err) + } + + if code, activated := getOrder(t, router, "pay001", 1); code != http.StatusOK || activated { + t.Fatalf("pending 不应开通:code=%d activated=%v", code, activated) + } + row, _ := st.GetForUser(ctx, 1, "pay001") + if row.Status != "created" { + t.Fatalf("pending 时本地台账应仍 created, got %q", row.Status) + } + var n int + _ = db.QueryRow(`SELECT COUNT(*) FROM subscriptions WHERE user_id=1`).Scan(&n) + if n != 0 { + t.Fatalf("pending 不得开通订阅,得 %d", n) + } +}