Files
wangjia e0c7baa7cb fix(sec): 安全评审加固——goroutine panic 恢复 + 下单 body 上限 + 错误脱敏
- 后台 goroutine (notifyBiz/SyncPending) 加 defer recover,防 panic 拖垮进程
- 下单请求体 64KB 上限 (MaxBytesReader),防超大 body 内存 DoS
- 下单错误不再回传原始 err.Error(),内部详情只进日志
- (另: pay nginx 加 per-IP 限流 limit_req,服务器侧配置)

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_019UQmqWmV67sXGLrb3U1XXn
2026-07-03 23:53:11 +08:00

409 lines
13 KiB
Go

package service
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net/http"
"strconv"
"strings"
"time"
"github.com/google/uuid"
"gorm.io/gorm"
"github.com/wangjia/pay/config"
"github.com/wangjia/pay/internal/channel"
"github.com/wangjia/pay/internal/model"
"github.com/wangjia/pay/internal/util"
)
var (
ErrProductNotFound = errors.New("套餐不存在或已下架")
ErrAmountMismatch = errors.New("回调金额与订单金额不符")
)
type OrderService struct {
db *gorm.DB
reg *channel.Registry
baseURL string
}
func NewOrderService(db *gorm.DB, reg *channel.Registry, baseURL string) *OrderService {
return &OrderService{db: db, reg: reg, baseURL: baseURL}
}
// BizParams 业务对接下单参数(独立收款场景全空)。
type BizParams struct {
System string // 业务来源,如 jiu
Ref string // 业务引用,如 jiu 的 purchase_id
ReturnURL string // 自定义付款后跳回地址(空则用 pay 默认结果页)
}
// prepare 校验套餐、取渠道、落库一张待支付订单(金额一律取服务端套餐价,不信任前端)。
func (s *OrderService) prepare(productID uint64, clientIP string, biz BizParams) (channel.Channel, *model.Merchant, *model.Order, error) {
var p model.Product
if err := s.db.First(&p, "id = ? AND active = ?", productID, true).Error; err != nil {
return nil, nil, nil, ErrProductNotFound
}
ch, m, err := s.reg.ByMerchantID(p.MerchantID)
if err != nil {
return nil, nil, nil, err
}
order := &model.Order{
OutTradeNo: util.NewOutTradeNo(m.Code),
MerchantID: m.ID,
Channel: m.Channel,
ProductID: p.ID,
Subject: p.Name,
Amount: p.Price, // 权威金额
Status: model.OrderPending,
ClientIP: clientIP,
BizSystem: biz.System,
BizRef: biz.Ref,
}
if err := s.db.Create(order).Error; err != nil {
return nil, nil, nil, fmt.Errorf("创建订单失败: %w", err)
}
return ch, m, order, nil
}
func (s *OrderService) notifyURL(channel string) string {
return s.baseURL + "/api/v1/notify/" + channel
}
// returnURL 付款后同步跳转地址:业务方传了自定义地址就用它(拼上 out_trade_no),否则用 pay 默认结果页。
func (s *OrderService) returnURL(custom, outTradeNo string) string {
if custom == "" {
return s.baseURL + "/result?out_trade_no=" + outTradeNo
}
sep := "?"
if strings.Contains(custom, "?") {
sep = "&"
}
return custom + sep + "out_trade_no=" + outTradeNo
}
// Create 网页支付下单,返回收银台跳转 URL。isMobile=true 时支付宝走手机网站支付(拉起 App)。
func (s *OrderService) Create(ctx context.Context, productID uint64, clientIP string, isMobile bool, biz BizParams) (string, *model.Order, error) {
ch, m, order, err := s.prepare(productID, clientIP, biz)
if err != nil {
return "", nil, err
}
payURL, err := ch.PagePay(ctx, channel.CreateReq{
OutTradeNo: order.OutTradeNo,
Subject: order.Subject,
Amount: order.Amount,
NotifyURL: s.notifyURL(m.Channel),
ReturnURL: s.returnURL(biz.ReturnURL, order.OutTradeNo),
IsMobile: isMobile,
})
if err != nil {
return "", nil, err
}
return payURL, order, nil
}
// CreateQR 扫码(当面付)下单,返回二维码码串供前端渲染。
func (s *OrderService) CreateQR(ctx context.Context, productID uint64, clientIP string) (string, *model.Order, error) {
ch, m, order, err := s.prepare(productID, clientIP, BizParams{})
if err != nil {
return "", nil, err
}
qr, err := ch.PreCreate(ctx, channel.CreateReq{
OutTradeNo: order.OutTradeNo,
Subject: order.Subject,
Amount: order.Amount,
NotifyURL: s.notifyURL(m.Channel),
})
if err != nil {
return "", nil, err
}
return qr, order, nil
}
// HandleAlipayNotify 处理支付宝异步回调:反查商户 → 验签 → 核对金额 → 幂等更新。
// 返回 nil 表示已正确处理(调用方应给支付宝回 "success")。
func (s *OrderService) HandleAlipayNotify(ctx context.Context, r *http.Request) error {
if err := r.ParseForm(); err != nil {
return fmt.Errorf("解析回调失败: %w", err)
}
appID := r.PostFormValue("app_id")
outTradeNo := r.PostFormValue("out_trade_no")
ch, m, err := s.reg.AlipayByAppID(appID)
if err != nil {
s.logNotify("alipay", outTradeNo, false, "not_found", r.Form.Encode())
return err
}
res, err := ch.VerifyNotify(ctx, r)
if err != nil {
s.logNotify("alipay", outTradeNo, false, "verify_failed", r.Form.Encode())
return err
}
result, err := s.applyPaid(m, res)
s.logNotify("alipay", res.OutTradeNo, true, result, res.Raw)
return err
}
// HandleWechatNotify 处理微信异步回调:解密验签 → 按 out_trade_no 定位订单/商户 → 核对金额 → 幂等更新。
// 微信回调正文加密且不含明文商户路由信息,故先用任一启用的微信商户凭证解密(同主体共用),
// 再按解密出的 out_trade_no 找到订单实际归属的商户入账。返回 nil 表示已正确处理。
func (s *OrderService) HandleWechatNotify(ctx context.Context, r *http.Request) error {
ch, _, err := s.reg.FirstWechat()
if err != nil {
s.logNotify("wechat", "", false, "no_merchant", "")
return err
}
res, err := ch.VerifyNotify(ctx, r)
if err != nil {
s.logNotify("wechat", "", false, "verify_failed", "")
return err
}
var o model.Order
if err := s.db.First(&o, "out_trade_no = ?", res.OutTradeNo).Error; err != nil {
s.logNotify("wechat", res.OutTradeNo, true, "not_found", res.Raw)
return fmt.Errorf("订单不存在: %s", res.OutTradeNo)
}
var m model.Merchant
if err := s.db.First(&m, o.MerchantID).Error; err != nil {
s.logNotify("wechat", res.OutTradeNo, true, "merchant_not_found", res.Raw)
return fmt.Errorf("订单 %s 的商户不存在: %w", res.OutTradeNo, err)
}
result, err := s.applyPaid(&m, res)
s.logNotify("wechat", res.OutTradeNo, true, result, res.Raw)
return err
}
// applyPaid 在一个事务里完成「金额核对 + 幂等置为已支付」。返回处理结果标记。
func (s *OrderService) applyPaid(m *model.Merchant, res *channel.NotifyResult) (string, error) {
if !res.Paid {
return "ignored", nil // 非成功状态(如 WAIT_BUYER_PAY),确认收到即可
}
var resultTag string
err := s.db.Transaction(func(tx *gorm.DB) error {
var o model.Order
if err := tx.First(&o, "out_trade_no = ? AND merchant_id = ?", res.OutTradeNo, m.ID).Error; err != nil {
resultTag = "not_found"
return fmt.Errorf("订单不存在: %s", res.OutTradeNo)
}
if o.Status == model.OrderPaid {
resultTag = "duplicate" // 幂等:已处理过,直接成功返回
return nil
}
if !util.AmountEqual(o.Amount, res.Amount) {
resultTag = "amount_mismatch"
return ErrAmountMismatch
}
now := time.Now()
upd := tx.Model(&model.Order{}).
Where("out_trade_no = ? AND status = ?", o.OutTradeNo, model.OrderPending).
Updates(map[string]any{
"status": model.OrderPaid,
"trade_no": res.TradeNo,
"buyer_logon_id": res.BuyerLogonID,
"paid_at": &now,
})
if upd.Error != nil {
return upd.Error
}
if upd.RowsAffected == 0 {
resultTag = "duplicate" // 并发下被另一路(如查单)先置位
return nil
}
resultTag = "processed"
log.Printf("[notify] 订单 %s 已支付 trade_no=%s amount=%s", o.OutTradeNo, res.TradeNo, res.Amount)
return nil
})
// 新入账成功且属业务对接单:异步回调业务系统 webhook(失败由后台重试兜底,不阻塞支付回调响应)。
if err == nil && resultTag == "processed" {
go s.notifyBizByOutTradeNo(res.OutTradeNo)
}
return resultTag, err
}
func (s *OrderService) logNotify(ch, outTradeNo string, verified bool, result, raw string) {
_ = s.db.Create(&model.NotifyLog{
Channel: ch,
OutTradeNo: outTradeNo,
Verified: verified,
Result: result,
Raw: raw,
}).Error
}
// GetByOutTradeNo 供前端结果页轮询。
func (s *OrderService) GetByOutTradeNo(outTradeNo string) (*model.Order, error) {
var o model.Order
if err := s.db.First(&o, "out_trade_no = ?", outTradeNo).Error; err != nil {
return nil, err
}
return &o, nil
}
// SyncPending 兜底:把近期待支付订单拿去主动查单,命中已支付则补记(防回调丢失)。
func (s *OrderService) SyncPending(ctx context.Context, maxAge time.Duration) {
// 后台循环调用;恢复 panic 防止拖垮整个进程。
defer func() {
if r := recover(); r != nil {
log.Printf("[query_sync] 发生 panic 已恢复: %v", r)
}
}()
var orders []model.Order
cutoff := time.Now().Add(-maxAge)
if err := s.db.Where("status = ? AND created_at > ?", model.OrderPending, cutoff).
Limit(100).Find(&orders).Error; err != nil {
log.Printf("[query_sync] 查询待支付订单失败: %v", err)
return
}
for i := range orders {
o := &orders[i]
ch, m, err := s.reg.ByMerchantID(o.MerchantID)
if err != nil {
continue
}
qr, err := ch.Query(ctx, o.OutTradeNo)
if err != nil || qr == nil || !qr.Found || !qr.Paid {
continue
}
result, _ := s.applyPaid(m, &channel.NotifyResult{
OutTradeNo: qr.OutTradeNo,
TradeNo: qr.TradeNo,
Amount: qr.Amount,
Paid: true,
})
if result == "processed" {
log.Printf("[query_sync] 订单 %s 经主动查单补记为已支付", o.OutTradeNo)
}
}
}
// StartQuerySync 启动后台查单兜底循环。
func (s *OrderService) StartQuerySync(interval, maxAge time.Duration) {
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for range ticker.C {
s.SyncPending(context.Background(), maxAge)
}
}()
}
// notifyBizByOutTradeNo 向业务系统推送「支付成功」webhook(签名)。成功则置 BizNotified,失败留待重试。
func (s *OrderService) notifyBizByOutTradeNo(outTradeNo string) {
// 该函数在 goroutine 中运行(applyPaid 异步触发 + 重试循环);Go 中未恢复的 goroutine panic 会整进程崩溃。
defer func() {
if r := recover(); r != nil {
log.Printf("[biz_notify] 订单 %s 回调发生 panic 已恢复: %v", outTradeNo, r)
}
}()
var o model.Order
if err := s.db.First(&o, "out_trade_no = ?", outTradeNo).Error; err != nil {
return
}
if o.BizSystem == "" || o.BizNotified || o.Status != model.OrderPaid {
return
}
cfg, ok := config.C.BizByName(o.BizSystem)
if !ok {
log.Printf("[biz_notify] 订单 %s 业务系统 %s 未配置回调,跳过", o.OutTradeNo, o.BizSystem)
return
}
var p model.Product
_ = s.db.First(&p, o.ProductID).Error // 取 biz_code;失败则为空
paidAt := ""
if o.PaidAt != nil {
paidAt = o.PaidAt.Format(time.RFC3339)
}
body, _ := json.Marshal(map[string]any{
"out_trade_no": o.OutTradeNo,
"biz_system": o.BizSystem,
"biz_ref": o.BizRef,
"product_biz_code": p.BizCode,
"amount": o.Amount,
"trade_no": o.TradeNo,
"channel": o.Channel,
"paid_at": paidAt,
})
ts := strconv.FormatInt(time.Now().Unix(), 10)
nonce := uuid.NewString()
sign := util.HMACSign(cfg.Secret, o.BizSystem, ts, nonce, string(body))
req, err := http.NewRequest(http.MethodPost, cfg.CallbackURL, bytes.NewReader(body))
if err != nil {
log.Printf("[biz_notify] 订单 %s 构造回调请求失败: %v", o.OutTradeNo, err)
return
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("X-Pay-System", o.BizSystem)
req.Header.Set("X-Pay-Timestamp", ts)
req.Header.Set("X-Pay-Nonce", nonce)
req.Header.Set("X-Pay-Sign", sign)
entry := &model.BizNotifyLog{OutTradeNo: o.OutTradeNo, BizSystem: o.BizSystem, URL: cfg.CallbackURL, Payload: string(body)}
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req)
success := false
if err != nil {
entry.RespBody = err.Error()
} else {
rb, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
resp.Body.Close()
entry.RespCode = resp.StatusCode
entry.RespBody = string(rb)
// 约定:业务方返回 HTTP 200 且响应含 SUCCESS 视为受理成功。
if resp.StatusCode == http.StatusOK && strings.Contains(strings.ToUpper(string(rb)), "SUCCESS") {
success = true
}
}
entry.OK = success
_ = s.db.Create(entry).Error
if success {
s.db.Model(&model.Order{}).Where("out_trade_no = ?", o.OutTradeNo).Update("biz_notified", true)
log.Printf("[biz_notify] 订单 %s 已成功回调 %s", o.OutTradeNo, o.BizSystem)
} else {
log.Printf("[biz_notify] 订单 %s 回调 %s 失败(code=%d),等待重试", o.OutTradeNo, o.BizSystem, entry.RespCode)
}
}
// NotifyBizPending 兜底重试:把已支付但未成功回调业务方的订单再推一次。
func (s *OrderService) NotifyBizPending() {
var orders []model.Order
cutoff := time.Now().Add(-24 * time.Hour)
if err := s.db.Where("status = ? AND biz_system <> '' AND biz_notified = ? AND created_at > ?",
model.OrderPaid, false, cutoff).Limit(50).Find(&orders).Error; err != nil {
log.Printf("[biz_notify] 查询待回调订单失败: %v", err)
return
}
for i := range orders {
s.notifyBizByOutTradeNo(orders[i].OutTradeNo)
}
}
// StartBizNotifyRetry 启动业务回调重试循环(兜底 webhook 丢失/业务方短暂不可用)。
func (s *OrderService) StartBizNotifyRetry(interval time.Duration) {
go func() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for range ticker.C {
s.NotifyBizPending()
}
}()
}