Files
wangjia a62a2b1797
ci-pangolin / Lint — shellcheck (pull_request) Successful in 11s
ci-pangolin / Redline Scan — 脱敏 (UI 文案) (pull_request) Successful in 25s
ci-pangolin / Cleartext Scan — Android 禁明文 (pull_request) Successful in 19s
ci-pangolin / OpenAPI Sync Check (pull_request) Successful in 41s
ci-pangolin / Portable SQL — 可移植性 (mysql/sqlite) (pull_request) Successful in 20s
ci-pangolin / Codegen Drift — token 生成物未漂移 (pull_request) Successful in 4s
ci-pangolin / DS-flow — 原型/跨端同源/代码色单源闸 (pull_request) Successful in 5s
ci-pangolin / Go — build + test (pull_request) Failing after 14s
ci-pangolin / E2E Smoke — L4 进程级端到端 (pull_request) Failing after 11s
ci-pangolin / Go — integration (mysql/redis testcontainers) (pull_request) Failing after 4m44s
ci-pangolin / Golden — 视觉回归 (全量:components/auth/desktop/tablet) (pull_request) Failing after 20s
ci-pangolin / Flutter — analyze + test (pull_request) Failing after 11m55s
fix(server/pay): promo 限购 settle 侧幂等复查 + 000027 部分唯一索引(TOCTOU)
CreateOrder 下单时的 HasPaidPurchase 只是裸 SELECT 无锁,并发/多挂起单可绕过
promo SKU「每账号限购一次」。两层修:
① webhook.go settle 在锁行、开通前对 item.Promo 的 SKU 复查一次(排除本单),
   命中说明另一笔同 user+SKU 订单已抢先 settle,跳过发放(不二次 +N 天)、
   仍 ack(否则 pay 无限重投)。新增 store.HasPaidPurchaseExcludingTx /
   MarkDuplicatePromoTx——重复单标记 canceled 而非 paid,避免自撞下面的
   唯一索引、也避免整笔 500 触发死循环重投。
② migration 000027(sqlite):部分唯一索引 ux_pay_promo_paid ON
   pay_purchases(user_id, sku) WHERE status='paid' AND sku='pro_month_promo',
   兜底防止任何路径把同一用户的 promo 单二次写成 paid。mysql 8 不支持部分
   索引,000027 mysql 侧是 no-op 占位(仅对齐编号),该场景 mysql 只靠①的
   应用层复查兜底——两库防线强度不同,已在迁移文件与代码注释中记录。

副作用:新增迁移把 sqlite 迁移顶点从 26 推到 27,同步更新
internal/store/sqlite_migrate_test.go 的版本断言,以及
internal/store/codes_lib_migrate_test.go 手动 Steps(-1) 序列(补一步跳过
000027,才能精确落在 000022 边界,这条测试是硬编码步数的)。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-13 15:53:56 +08:00

230 lines
8.5 KiB
Go

package pay
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"strconv"
"time"
"github.com/redis/go-redis/v9"
"github.com/wangjia/pangolin/server/internal/codes"
)
// Granter 抽象 codes.Service 的支付授予入口(测试可替身;生产传 *codes.Service)。
type Granter interface {
GrantPaidSubscriptionTx(ctx context.Context, tx *sql.Tx, userID int64, plan codes.PlanCode, days int, ref string) (int64, time.Time, error)
}
// WebhookHandler 接收 pay 的 payment.succeeded 出站 webhook。
// 验签与 pay verifyBizSign 对称:同 secret,parts=[system, ts, nonce, rawBody],
// ±tolerance 时间窗。幂等三层:nonce SETNX(传输重放)→ out_trade_no 锁内
// CAS(业务幂等,重投唯一可靠键)→ biz_ref 兜底(台账缺行自修复)。
type WebhookHandler struct {
store *Store
granter Granter
db *sql.DB
rdb *redis.Client // 可为 nil:跳过 nonce 层,业务幂等仍成立
system string
secret string
tolerance time.Duration
nonceTTL time.Duration
now func() time.Time // 测试注入
rewarder Rewarder // 可选:首充邀请奖励钩子(nil 则跳过)
noticer Noticer // 可选:购买开通到账通知钩子(nil 则跳过)
}
func NewWebhookHandler(store *Store, granter Granter, db *sql.DB, rdb *redis.Client,
system, secret string, tolerance, nonceTTL time.Duration) *WebhookHandler {
return &WebhookHandler{store: store, granter: granter, db: db, rdb: rdb,
system: system, secret: secret, tolerance: tolerance, nonceTTL: nonceTTL, now: time.Now}
}
// Rewarder(可选)在首充同事务内发放邀请首充奖励;nil 则跳过。
type Rewarder interface {
OnFirstPaidTx(ctx context.Context, tx *sql.Tx, inviteeID int64, now time.Time) error
}
// SetRewarder 挂载首充邀请奖励钩子(reward.Service 满足此接口)。
func (h *WebhookHandler) SetRewarder(r Rewarder) { h.rewarder = r }
// Noticer 抽象 notices.Store 的事务内插入入口(测试可替身;生产传 notices.NewStore(db))。
type Noticer interface {
InsertNoticeTx(ctx context.Context, tx *sql.Tx, userID int64, typ, titleZH, titleEN, bodyZH, bodyEN, link string, now time.Time) error
}
// SetNoticer 挂载购买开通到账通知钩子;为 nil 时跳过(装配前兼容)。
func (h *WebhookHandler) SetNoticer(n Noticer) { h.noticer = n }
// webhookEvent 对应 pay settle.go::enqueuePaymentSucceeded 的 payload
// (注意:payment.succeeded 无 refund_id 字段)。
type webhookEvent struct {
EventType string `json:"event_type"`
OutTradeNo string `json:"out_trade_no"`
BizSystem string `json:"biz_system"`
BizRef string `json:"biz_ref"`
ProductBizCode string `json:"product_biz_code"`
AmountMinor int64 `json:"amount_minor"`
Currency string `json:"currency"`
Channel string `json:"channel"`
PaidAt string `json:"paid_at"` // RFC3339
}
func (h *WebhookHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
body, err := io.ReadAll(io.LimitReader(r.Body, 64<<10))
if err != nil {
http.Error(w, "read body", http.StatusBadRequest)
return
}
if !h.verify(r, body) {
http.Error(w, "signature verification failed", http.StatusUnauthorized)
return
}
// nonce 防重放(仅传输层;pay 每次重投换新 nonce,业务幂等靠 out_trade_no)。
if h.rdb != nil {
if nonce := r.Header.Get("X-Pay-Nonce"); nonce != "" {
ok, err := h.rdb.SetNX(r.Context(), "pay:webhook:nonce:"+nonce, 1, h.nonceTTL).Result()
if err == nil && !ok {
writeSuccess(w) // 同 nonce 重放:已处理过,直接确认
return
}
}
}
var ev webhookEvent
if err := json.Unmarshal(body, &ev); err != nil {
http.Error(w, "bad payload", http.StatusBadRequest)
return
}
if ev.EventType != "payment.succeeded" {
// 事件白名单外(pay 侧只应配 payment.succeeded):确认不处理,免重投。
slog.Warn("pay webhook: 未知事件已 ack 未处理", "event_type", ev.EventType, "out_trade_no", ev.OutTradeNo)
writeSuccess(w)
return
}
if err := h.settle(r.Context(), &ev); err != nil {
slog.Error("pay webhook 开通失败(pay 将退避重投)", "order_no", ev.OutTradeNo, "err", err)
http.Error(w, "settle failed", http.StatusInternalServerError)
return
}
writeSuccess(w)
}
// writeSuccess:pay 的 ACK 判据是 HTTP 200 且 body 含 "SUCCESS"(大写包含)。
func writeSuccess(w http.ResponseWriter) {
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("SUCCESS"))
}
func (h *WebhookHandler) verify(r *http.Request, body []byte) bool {
if r.Header.Get("X-Pay-System") != h.system {
return false
}
ts := r.Header.Get("X-Pay-Timestamp")
nonce := r.Header.Get("X-Pay-Nonce")
sign := r.Header.Get("X-Pay-Sign")
if ts == "" || nonce == "" || sign == "" {
return false
}
tsi, err := strconv.ParseInt(ts, 10, 64)
if err != nil {
return false
}
tol := int64(h.tolerance.Seconds())
if d := h.now().Unix() - tsi; d > tol || d < -tol {
return false
}
return hmacVerify(h.secret, sign, h.system, ts, nonce, string(body))
}
// settle 幂等开通:锁台账行 → created→paid 翻转 + 同事务 grant(叠加语义
// 复用 codes.applySubscription)。已 paid 直接返回 nil(重投/并发输家)。
// canceled 行也照常开通——钱已实收,本地 cancel 只是未支付单的整理。
func (h *WebhookHandler) settle(ctx context.Context, ev *webhookEvent) error {
item, ok := CatalogBySKU(ev.ProductBizCode)
if !ok {
return fmt.Errorf("未知 product_biz_code %q(与 pay 种子漂移?)", ev.ProductBizCode)
}
tx, err := h.store.BeginTx(ctx)
if err != nil {
return err
}
defer func() { _ = tx.Rollback() }()
var purchaseID, userID int64
row, err := h.store.LockByOutTradeNoTx(ctx, tx, ev.OutTradeNo)
switch {
case err == sql.ErrNoRows:
// 台账缺行(下单后本地写失败)→ 按 biz_ref=用户 uuid 兜底定位补建。
// 与正常路径(row.UserID 直接开通,不看 user 状态)对齐:钱已实收,
// 兜底定位不应因用户被停用(suspended)而拒绝开通——否则静默丢钱。
if err := h.db.QueryRowContext(ctx,
`SELECT id FROM users WHERE uuid = ?`, ev.BizRef).Scan(&userID); err != nil {
return fmt.Errorf("biz_ref %q 定位用户失败: %w", ev.BizRef, err)
}
purchaseID, err = h.store.InsertFromWebhookTx(ctx, tx, userID, ev.BizRef, ev.ProductBizCode, ev.OutTradeNo, ev.Channel)
if err != nil {
return err
}
case err != nil:
return err
case row.Status == "paid":
return nil // 幂等:已消费,直接 SUCCESS
default:
purchaseID, userID = row.ID, row.UserID
}
paidAt, perr := time.Parse(time.RFC3339, ev.PaidAt)
if perr != nil {
paidAt = h.now().UTC()
}
// C2 安全修复(promo 限购 TOCTOU):CreateOrder 下单时的 HasPaidPurchase 只在
// 下单一刻查、裸 SELECT 无锁——并发/多笔挂起单可绕过"每账号限购一次"。这里
// 在锁行、开通前对 Promo SKU 再复查一次:命中说明另一笔同 user+SKU 的订单
// 已抢先 settle 为 paid,本单跳过发放(不二次 +N 天),仍需 ACK 200,否则 pay
// 会把 500 当失败无限重投。sqlite 侧另有 migration 000027 的部分唯一索引
// (user_id, sku) WHERE status='paid' 兜底;mysql 不支持部分索引,本检查是
// mysql 侧唯一防线(见该迁移文件注释)。
if item.Promo {
dup, err := h.store.HasPaidPurchaseExcludingTx(ctx, tx, userID, ev.ProductBizCode, purchaseID)
if err != nil {
return err
}
if dup {
slog.Warn("pay webhook: promo 限购 TOCTOU 命中,跳过发放并吞掉重复单",
"order_no", ev.OutTradeNo, "user_id", userID, "sku", ev.ProductBizCode)
if err := h.store.MarkDuplicatePromoTx(ctx, tx, purchaseID, ev.AmountMinor, ev.Currency, ev.Channel, paidAt); err != nil {
return err
}
return tx.Commit()
}
}
subID, _, err := h.granter.GrantPaidSubscriptionTx(ctx, tx, userID,
codes.PlanCode(item.Plan), item.Days, "pay:"+ev.OutTradeNo)
if err != nil {
return err
}
if err := h.store.MarkPaidTx(ctx, tx, purchaseID, ev.AmountMinor, ev.Currency, ev.Channel, subID, paidAt); err != nil {
return err
}
if h.rewarder != nil {
if err := h.rewarder.OnFirstPaidTx(ctx, tx, userID, h.now().UTC()); err != nil {
return err // 同事务:发奖失败则整笔回滚,webhook 重试
}
}
if h.noticer != nil {
zh := fmt.Sprintf("已开通 Pro · %d 天", item.Days)
en := fmt.Sprintf("Pro activated · %d days", item.Days)
if err := h.noticer.InsertNoticeTx(ctx, tx, userID, "reward", zh, en, "", "", "", h.now().UTC()); err != nil {
return err // 同事务:通知插入失败则整笔回滚,webhook 重试
}
}
return tx.Commit()
}