614 lines
22 KiB
Go
614 lines
22 KiB
Go
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)
|
||
}
|
||
}
|
||
}
|
||
}()
|
||
}
|