Files
2026-09-15 12:46:29 +08:00

614 lines
22 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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)
}
}
}
}()
}