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) } } } }() }