增加支付功能

This commit is contained in:
2026-09-15 12:46:29 +08:00
parent 943c3708b0
commit a82b300b1a
46 changed files with 11675 additions and 21 deletions
+613
View File
@@ -0,0 +1,613 @@
package payment
import (
"context"
"errors"
"fmt"
"net/http"
"strings"
"time"
"server/models"
"github.com/beego/beego/v2/client/orm"
beelog "github.com/beego/beego/v2/core/logs"
)
/* =============================================================
* 支付服务核心:创建支付单 / 查询 / 状态机 / 回调幂等 / 兜底退回 / 补偿任务
*
* 与业务模块的解耦点:
* 1. 业务订单的「已支付」同步通过 RegisterOrderSyncHook 注册钩子完成;
* 未注册钩子时视为无需同步(order_synced 直接置 1)。
* 2. 佣金推广方:下单入参可携带(快照到支付单),未携带时由
* RegisterPromoterResolver 注册的钩子推导(如按租户归属的渠道伙伴)。
* ============================================================= */
// OrderSyncFunc 业务订单状态同步钩子
type OrderSyncFunc func(ctx context.Context, order *models.PlatformPaymentOrder) error
// PromoterResolveFunc 推广方解析钩子
type PromoterResolveFunc func(ctx context.Context, order *models.PlatformPaymentOrder) (promoterID, promoterName, promoterType string, ok bool)
var (
orderSyncHook OrderSyncFunc
promoterResolve PromoterResolveFunc
)
// RegisterOrderSyncHook 注册业务订单同步钩子(在业务模块 init 时调用)
func RegisterOrderSyncHook(fn OrderSyncFunc) { orderSyncHook = fn }
// RegisterPromoterResolver 注册推广方解析钩子
func RegisterPromoterResolver(fn PromoterResolveFunc) { promoterResolve = fn }
// CreateInput 创建支付单入参(租户端下单页 -> POST /backend/payment/create)
type CreateInput struct {
OutTradeNo string // 业务订单号(必填)
OrderType string // platform_usage / module_shop / service_fee
TenantID string
TenantName string
Amount int64 // 分
Channel string // 渠道标识(必填)
Subject string
ReturnURL string // 支付完成同步跳回地址
ClientIP string
PayType string // qr / web / h5 / jsapi;空值由适配器给默认
OpenID string // JSAPI 必填
PromoterID string
PromoterName string
PromoterType string
ExpireMinutes int // 支付有效期(分钟),默认 30
}
const defaultExpireMinutes = 30
// GetOrderByPayNo 按支付单号查询
func GetOrderByPayNo(payNo string) (*models.PlatformPaymentOrder, error) {
row := &models.PlatformPaymentOrder{}
err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("pay_no", payNo).
Filter("delete_time__isnull", true).
One(row)
if err != nil {
return nil, err
}
return row, nil
}
// LogStateChange 记录状态流转流水
func LogStateChange(payNo, fromStatus, toStatus, operator, remark string) {
row := &models.PlatformPaymentOrderLog{
PayNo: payNo,
FromStatus: fromStatus,
ToStatus: toStatus,
Operator: operator,
Remark: remark,
}
if _, err := models.Orm.Insert(row); err != nil {
beelog.Warn("支付单流转日志写入失败: %s %v", payNo, err)
}
}
// CreatePayment 创建支付单并调用渠道下单
// 幂等:同一 out_trade_no + channel 存在未终态支付单时,直接复用并重新拉起支付参数。
func CreatePayment(ctx context.Context, in CreateInput) (*models.PlatformPaymentOrder, *PayParams, error) {
if strings.TrimSpace(in.OutTradeNo) == "" {
return nil, nil, errors.New("业务订单号不能为空")
}
if in.Amount <= 0 {
return nil, nil, errors.New("支付金额必须大于 0")
}
adapter, err := GetChannelAdapter(in.Channel)
if err != nil {
return nil, nil, err
}
cfg, err := LoadEnabledChannelConfig(in.Channel)
if err != nil {
return nil, nil, err
}
// 幂等:复用未终态支付单
var exist models.PlatformPaymentOrder
err = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("out_trade_no", in.OutTradeNo).
Filter("channel", in.Channel).
Filter("delete_time__isnull", true).
Filter("status__in", models.PayStatusCreated, models.PayStatusPending, models.PayStatusPaying).
One(&exist)
if err == nil {
params, perr := adapter.Prepay(ctx, &exist, cfg, PrepayOption{PayType: in.PayType, OpenID: in.OpenID})
if perr != nil {
return nil, nil, perr
}
return &exist, params, nil
} else if err != orm.ErrNoRows {
return nil, nil, fmt.Errorf("查询支付单失败: %w", err)
}
payNo := NextPayNo()
now := time.Now()
minutes := in.ExpireMinutes
if minutes <= 0 {
minutes = defaultExpireMinutes
}
expireAt := now.Add(time.Duration(minutes) * time.Minute)
order := &models.PlatformPaymentOrder{
PayNo: payNo,
OutTradeNo: in.OutTradeNo,
OrderType: in.OrderType,
Subject: in.Subject,
TenantID: in.TenantID,
TenantName: in.TenantName,
Amount: in.Amount,
OrderAmount: in.Amount,
Currency: "CNY",
Channel: in.Channel,
MerchantNo: cfg.MerchantNo,
Status: models.PayStatusCreated,
ClientIP: in.ClientIP,
ReturnURL: in.ReturnURL,
NotifyURL: cfg.CallbackURL,
ExpireAt: &expireAt,
PromoterID: in.PromoterID,
PromoterName: in.PromoterName,
PromoterType: in.PromoterType,
}
// 下单未携带推广方时,交给业务模块推导(如租户的签约渠道伙伴)
if order.PromoterID == "" && promoterResolve != nil {
if pid, pname, ptype, ok := promoterResolve(ctx, order); ok {
order.PromoterID, order.PromoterName, order.PromoterType = pid, pname, ptype
}
}
if _, err := models.Orm.Insert(order); err != nil {
return nil, nil, fmt.Errorf("创建支付单失败: %w", err)
}
LogStateChange(order.PayNo, "", models.PayStatusCreated, "租户端下单页", "创建支付单,渠道:"+cfg.Name)
params, err := adapter.Prepay(ctx, order, cfg, PrepayOption{PayType: in.PayType, OpenID: in.OpenID})
if err != nil {
// 渠道下单失败:保留支付单并置为 failed,便于排查与重下
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"status": models.PayStatusFailed, "update_time": now})
LogStateChange(order.PayNo, models.PayStatusCreated, models.PayStatusFailed, "支付服务", "渠道下单失败:"+err.Error())
return nil, nil, err
}
LogStateChange(order.PayNo, models.PayStatusCreated, models.PayStatusCreated, "支付服务", "渠道下单成功,返回支付参数")
if params.ChannelTradeNo != "" {
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"channel_trade_no": params.ChannelTradeNo, "update_time": time.Now()})
order.ChannelTradeNo = params.ChannelTradeNo
}
return order, params, nil
}
// QueryPayment 查询支付单;syncChannel=true 时主动向渠道查询并同步状态(兜底回调丢失)
func QueryPayment(ctx context.Context, payNo string, syncChannel bool) (*models.PlatformPaymentOrder, error) {
order, err := GetOrderByPayNo(payNo)
if err != nil {
return nil, fmt.Errorf("支付单不存在")
}
now := time.Now()
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"last_query_at": now, "update_time": now})
order.LastQueryAt = &now
if !syncChannel {
return order, nil
}
adapter, aerr := GetChannelAdapter(order.Channel)
cfg, cerr := LoadChannelConfig(order.Channel)
if aerr != nil || cerr != nil {
return order, nil
}
st, qerr := adapter.Query(ctx, order, cfg)
if qerr != nil {
LogStateChange(order.PayNo, order.Status, order.Status, "手动查询", "查询渠道状态失败:"+qerr.Error())
return order, nil
}
if _, serr := applyChannelState(ctx, order, st, "手动查询", "主动查询渠道状态同步"); serr != nil {
beelog.Warn("支付单 %s 渠道状态同步失败: %v", order.PayNo, serr)
}
return GetOrderByPayNo(payNo)
}
// HandleNotify 渠道异步通知统一入口:验签 -> 幂等 -> 状态机 -> 佣金,返回需回给渠道的响应体
func HandleNotify(ctx context.Context, channel string, r *http.Request, ip string) (string, error) {
adapter, err := GetChannelAdapter(channel)
if err != nil {
return "", err
}
cfg, err := LoadChannelConfig(channel) // 停用渠道的历史回调仍需处理
if err != nil {
return "", err
}
result, err := adapter.ParseNotify(ctx, r, cfg)
if err != nil {
insertCallbackLog(channel, "", "", "", 0, 0, false, false, err.Error(), "", ip)
return "", err
}
// 幂等:同一事件(channel + event_id)已成功处理过则直接命中
duplicate := existsCallbackHandled(channel, result.EventID, result.ChannelTradeNo, result.EventType)
insertCallbackLog(channel, result.PayNo, result.OutTradeNo, result.ChannelTradeNo, result.Amount, 1, duplicate, false, "", result.Raw, ip)
if duplicate {
return result.AckBody, nil
}
order, oerr := GetOrderByPayNo(result.PayNo)
if oerr != nil {
// 再尝试按业务订单号找最近一笔(回调里可能只回传 out_trade_no)
order, oerr = findLatestOrderByOutTradeNo(result.OutTradeNo)
}
if oerr != nil {
insertCallbackLog(channel, result.PayNo, result.OutTradeNo, result.ChannelTradeNo, result.Amount, 1, false, false, "支付单不存在", result.Raw, ip)
return "", fmt.Errorf("支付单不存在: %s", result.PayNo)
}
// 非成功/失败/关闭事件只记录日志
if result.TradeState != StateSuccess && result.TradeState != StateFailed && result.TradeState != StateClosed {
insertCallbackLog(channel, result.PayNo, result.OutTradeNo, result.ChannelTradeNo, result.Amount, 1, false, true, "事件无需变更状态: "+result.EventType, result.Raw, ip)
return result.AckBody, nil
}
changed, serr := applyChannelState(ctx, order, &ChannelState{
ChannelTradeNo: result.ChannelTradeNo,
TradeState: result.TradeState,
Amount: result.Amount,
}, "渠道回调", "渠道通知同步,事件:"+result.EventType)
if serr != nil {
insertCallbackLog(channel, result.PayNo, result.OutTradeNo, result.ChannelTradeNo, result.Amount, 1, false, false, "状态更新失败: "+serr.Error(), result.Raw, ip)
return "", serr
}
if !changed {
insertCallbackLog(channel, result.PayNo, result.OutTradeNo, result.ChannelTradeNo, result.Amount, 1, false, true, "状态未变化(幂等)", result.Raw, ip)
} else {
insertCallbackLog(channel, result.PayNo, result.OutTradeNo, result.ChannelTradeNo, result.Amount, 1, false, true, "状态已同步为 "+result.TradeState, result.Raw, ip)
}
return result.AckBody, nil
}
// existsCallbackHandled 是否已存在处理成功的同事件回调
func existsCallbackHandled(channel, eventID, channelTradeNo, eventType string) bool {
qs := models.Orm.QueryTable(new(models.PlatformPaymentCallbackLog)).
Filter("channel", channel).
Filter("handle_result", 1).
Filter("delete_time__isnull", true)
if eventID != "" {
qs = qs.Filter("event_id", eventID)
} else {
qs = qs.Filter("channel_trade_no", channelTradeNo).Filter("event_type", eventType)
}
cnt, _ := qs.Count()
return cnt > 0
}
func insertCallbackLog(channel, payNo, outTradeNo, channelTradeNo string, amount int64, verifyResult int8, duplicate, handled bool, msg, raw, ip string) {
row := &models.PlatformPaymentCallbackLog{
Channel: channel,
PayNo: payNo,
OutTradeNo: outTradeNo,
ChannelTradeNo: channelTradeNo,
Amount: amount,
VerifyResult: verifyResult,
HandleMsg: msg,
ClientIP: ip,
}
if duplicate {
row.IsDuplicate = 1
}
if handled {
row.HandleResult = 1
}
if raw != "" {
raw = truncateStr(raw, 60000)
row.RawBody = &raw
}
if _, err := models.Orm.Insert(row); err != nil {
beelog.Warn("回调日志写入失败: %v", err)
}
}
func findLatestOrderByOutTradeNo(outTradeNo string) (*models.PlatformPaymentOrder, error) {
row := &models.PlatformPaymentOrder{}
err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("out_trade_no", outTradeNo).
Filter("delete_time__isnull", true).
OrderBy("-id").
One(row)
return row, err
}
// applyChannelState 把渠道状态落到本地(回调与主动查询共用),返回是否发生状态变更。
// 支付成功采用条件更新(仅 created/pending/paying 可置为 paid)保证并发幂等。
func applyChannelState(ctx context.Context, order *models.PlatformPaymentOrder, st *ChannelState, operator, remark string) (bool, error) {
if st == nil {
return false, nil
}
now := time.Now()
active := []interface{}{models.PayStatusCreated, models.PayStatusPending, models.PayStatusPaying}
switch st.TradeState {
case StateSuccess:
if order.Status == models.PayStatusPaid {
if order.ChannelTradeNo == "" && st.ChannelTradeNo != "" {
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"channel_trade_no": st.ChannelTradeNo, "update_time": now})
}
return false, nil
}
note := remark
if st.Amount > 0 && st.Amount != order.Amount {
note = fmt.Sprintf("%s(注意:渠道金额 %d 分与本地 %d 分不一致,请核对)", remark, st.Amount, order.Amount)
}
updates := orm.Params{
"status": models.PayStatusPaid,
"paid_at": now,
"notify_at": now,
"order_synced": 0,
"update_time": now,
}
if st.ChannelTradeNo != "" {
updates["channel_trade_no"] = st.ChannelTradeNo
}
n, err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Filter("status__in", active...).
Update(updates)
if err != nil {
return false, err
}
if n == 0 {
return false, nil // 并发下已被其他通知处理
}
order.Status = models.PayStatusPaid
order.PaidAt = &now
if st.ChannelTradeNo != "" {
order.ChannelTradeNo = st.ChannelTradeNo
}
LogStateChange(order.PayNo, models.PayStatusPaying, models.PayStatusPaid, operator, note)
// 1. 同步业务订单状态(钩子;失败保留 order_synced=0 由补偿任务重试)
if orderSyncHook == nil {
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"order_synced": 1, "update_time": time.Now()})
} else if serr := orderSyncHook(ctx, order); serr != nil {
LogStateChange(order.PayNo, models.PayStatusPaid, models.PayStatusPaid, operator, "业务订单同步失败:"+serr.Error())
} else {
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"order_synced": 1, "update_time": time.Now()})
}
// 2. 生成推广佣金(平台自有资金支出,失败不影响收款状态)
if _, cerr := OnOrderPaid(ctx, order); cerr != nil {
beelog.Warn("支付单 %s 佣金生成失败: %v", order.PayNo, cerr)
}
return true, nil
case StateClosed:
if order.Status == models.PayStatusPaid {
return false, nil
}
_, err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Filter("status__in", active...).
Update(orm.Params{"status": models.PayStatusClosed, "closed_at": now, "update_time": now})
if err != nil {
return false, err
}
order.Status = models.PayStatusClosed
LogStateChange(order.PayNo, models.PayStatusPending, models.PayStatusClosed, operator, remark)
return true, nil
case StateFailed:
if order.Status == models.PayStatusPaid {
return false, nil
}
_, err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Filter("status__in", active...).
Update(orm.Params{"status": models.PayStatusFailed, "update_time": now})
if err != nil {
return false, err
}
order.Status = models.PayStatusFailed
LogStateChange(order.PayNo, models.PayStatusPending, models.PayStatusFailed, operator, remark)
return true, nil
case StatePending:
if order.Status == models.PayStatusCreated {
_, err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"status": models.PayStatusPending, "update_time": now})
if err != nil {
return false, err
}
order.Status = models.PayStatusPending
LogStateChange(order.PayNo, models.PayStatusCreated, models.PayStatusPending, operator, remark)
return true, nil
}
return false, nil
default:
return false, nil
}
}
// RefundInput 手动原路退回入参(兜底能力)
type RefundInput struct {
PayNo string
Amount int64 // 分;<=0 表示全额退回
Reason string // 必填
OperatorID string
OperatorName string
}
// ManualRefund 手动原路退回(不做业务层退款流程)
func ManualRefund(ctx context.Context, in RefundInput) (*models.PlatformPaymentRefund, error) {
if strings.TrimSpace(in.Reason) == "" {
return nil, errors.New("退回原因不能为空")
}
order, err := GetOrderByPayNo(in.PayNo)
if err != nil {
return nil, fmt.Errorf("支付单不存在")
}
if order.Status != models.PayStatusPaid {
return nil, fmt.Errorf("仅支付成功的支付单支持原路退回(当前状态 %s)", order.Status)
}
remain := order.Amount - order.RefundAmount
amount := in.Amount
if amount <= 0 {
amount = remain
}
if amount > remain {
return nil, fmt.Errorf("退回金额 %d 分超过可退余额 %d 分", amount, remain)
}
refund := &models.PlatformPaymentRefund{
RefundNo: NextRefundNo(),
PayNo: order.PayNo,
OutTradeNo: order.OutTradeNo,
Channel: order.Channel,
Amount: amount,
Reason: in.Reason,
Status: models.RefundStatusProcessing,
OperatorID: in.OperatorID,
OperatorName: in.OperatorName,
}
if _, err := models.Orm.Insert(refund); err != nil {
return nil, fmt.Errorf("创建退回单失败: %w", err)
}
adapter, aerr := GetChannelAdapter(order.Channel)
cfg, cerr := LoadChannelConfig(order.Channel)
if aerr != nil || cerr != nil {
failRefund(refund, "渠道适配器/配置不可用")
return refund, fmt.Errorf("渠道适配器/配置不可用")
}
channelRefundNo, rerr := adapter.Refund(ctx, order, refund.RefundNo, amount, in.Reason, cfg)
now := time.Now()
if rerr != nil {
failRefund(refund, rerr.Error())
LogStateChange(order.PayNo, order.Status, order.Status, "平台管理员", "手动原路退回失败:"+rerr.Error())
return refund, rerr
}
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentRefund)).
Filter("id", refund.ID).
Update(orm.Params{"status": models.RefundStatusSuccess, "channel_refund_no": channelRefundNo, "finished_at": now, "update_time": now})
refund.Status = models.RefundStatusSuccess
refund.ChannelRefundNo = channelRefundNo
refund.FinishedAt = &now
newRefunded := order.RefundAmount + amount
toStatus := models.PayStatusRefunding
if newRefunded >= order.Amount {
toStatus = models.PayStatusRefunded
}
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", order.ID).
Update(orm.Params{"refund_amount": newRefunded, "status": toStatus, "update_time": now})
order.RefundAmount = newRefunded
order.Status = toStatus
LogStateChange(order.PayNo, models.PayStatusPaid, toStatus, "平台管理员", fmt.Sprintf("手动原路退回 ¥%s:%s", FenToYuan(amount), in.Reason))
return refund, nil
}
func failRefund(refund *models.PlatformPaymentRefund, msg string) {
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentRefund)).
Filter("id", refund.ID).
Update(orm.Params{"status": models.RefundStatusFailed, "fail_reason": truncateStr(msg, 480), "update_time": time.Now()})
refund.Status = models.RefundStatusFailed
refund.FailReason = msg
}
// CloseExpiredOrders 关闭超时未支付订单(补偿任务),返回关闭数量
func CloseExpiredOrders(ctx context.Context) (int, error) {
var orders []models.PlatformPaymentOrder
_, err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("delete_time__isnull", true).
Filter("status__in", models.PayStatusCreated, models.PayStatusPending, models.PayStatusPaying).
Filter("expire_at__lt", time.Now()).
All(&orders)
if err != nil {
return 0, err
}
closed := 0
now := time.Now()
for i := range orders {
_, err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", orders[i].ID).
Filter("status__in", models.PayStatusCreated, models.PayStatusPending, models.PayStatusPaying).
Update(orm.Params{"status": models.PayStatusClosed, "closed_at": now, "update_time": now})
if err != nil {
continue
}
LogStateChange(orders[i].PayNo, orders[i].Status, models.PayStatusClosed, "定时补偿", "支付超时自动关闭")
closed++
}
return closed, nil
}
// SyncPendingOrders 重试同步「已支付但业务订单未更新」的支付单(补偿任务),返回处理数量
func SyncPendingOrders(ctx context.Context) (int, error) {
if orderSyncHook == nil {
return 0, nil
}
var orders []models.PlatformPaymentOrder
_, err := models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("delete_time__isnull", true).
Filter("status", models.PayStatusPaid).
Filter("order_synced", 0).
All(&orders)
if err != nil {
return 0, err
}
done := 0
for i := range orders {
if serr := orderSyncHook(ctx, &orders[i]); serr != nil {
beelog.Warn("支付单 %s 业务订单同步重试失败: %v", orders[i].PayNo, serr)
continue
}
_, _ = models.Orm.QueryTable(new(models.PlatformPaymentOrder)).
Filter("id", orders[i].ID).
Update(orm.Params{"order_synced": 1, "update_time": time.Now()})
LogStateChange(orders[i].PayNo, models.PayStatusPaid, models.PayStatusPaid, "定时补偿", "业务订单同步成功(重试)")
done++
}
return done, nil
}
// StartPaymentScheduler 支付补偿任务:关闭超时订单 + 重试业务订单同步(5 分钟一轮)
func StartPaymentScheduler(stop <-chan struct{}) {
go func() {
ticker := time.NewTicker(5 * time.Minute)
defer ticker.Stop()
for {
select {
case <-stop:
return
case <-ticker.C:
ctx := context.Background()
if n, err := CloseExpiredOrders(ctx); err != nil {
beelog.Warn("支付超时关单失败: %v", err)
} else if n > 0 {
beelog.Info("支付超时关单 %d 笔", n)
}
if n, err := SyncPendingOrders(ctx); err != nil {
beelog.Warn("业务订单同步重试失败: %v", err)
} else if n > 0 {
beelog.Info("业务订单同步重试完成 %d 笔", n)
}
}
}
}()
}