feat(billing): 支付 P5.3 —— 掉单补偿定时器 + admin 订单流 + 日终对账
支付线封口。此前 pending 单只在「用户开着账单页轮询」时才查单确认——用户扫完码 关页面,钱付了、积分永不到账。 - 掉单补偿定时器(payment_reconcile.go):gateway 内每分钟扫 pending 微信单, 逐单 reconcileOrder 主动查单落态。把「用户在不在场」从入账链路摘掉。 reconcileOrder 从 BillingOrderStatus 抽出、前端轮询与定时器共用一份幂等 落态逻辑(不重蹈 GenerateReport/SubmitTask 的漂移)。渠道未配置时空转不炸。 - admin 订单流 GET /admin/orders(状态计数+全平台订单,可筛)。 - 日终对账 GET /admin/orders/reconcile:paid 单 ↔ 账本 grant 分录逐单比对, 抓 order_without_ledger(钱到了积分没给,最严重)/ ledger_without_paid_order。 - admin 计费页「充值订单与对账」块:计数卡片+订单流+一键对账。 ⚠️ live 抓到并修掉一个真 bug:OrderStats 复用同一个 gorm.DB 链式 Count 三次, WHERE 累加成 status=A AND status=B → 恒 0(订单流显示 2 单但计数全 0)。 改成每次起新 query builder。—— 又一次只有 live 才暴露的。 验证:go 6 包测试+tsc+41 vitest 全绿;live 造差异单对账正确抓出 order_without_ledger、清账后回零差异;补偿器启动日志+渠道未配置空转不炸; 浏览器验订单流卡片+一键对账绿条。TTL 过期路径需真渠道触发,部署后自然覆盖。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
@@ -93,27 +94,40 @@ func (h *Handler) BillingOrderStatus(c *gin.Context) {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "订单不存在"})
|
||||
return
|
||||
}
|
||||
if wc := h.pay.Current(); o.Status == store.OrderPending && wc != nil {
|
||||
if r, err := wc.QueryOrder(ctx, o.ID); err == nil {
|
||||
switch {
|
||||
case r.Paid && r.AmountFen == o.AmountFen:
|
||||
if _, err := h.db.MarkOrderPaid(ctx, o.ID, r.ChannelTxn); err == nil {
|
||||
o, _ = h.db.GetOrder(ctx, o.ID)
|
||||
}
|
||||
case r.Paid: // 金额对不上:不入账,人工对账(比错账便宜)
|
||||
c.JSON(http.StatusOK, gin.H{"order": o, "warn": "支付金额与订单不符,已挂起待人工核对"})
|
||||
return
|
||||
case r.Closed:
|
||||
_ = h.db.ExpireOrder(ctx, o.ID)
|
||||
updated, mismatch := h.reconcileOrder(ctx, o)
|
||||
if mismatch {
|
||||
c.JSON(http.StatusOK, gin.H{"order": updated, "warn": "支付金额与订单不符,已挂起待人工核对"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"order": updated})
|
||||
}
|
||||
|
||||
// reconcileOrder 对一张 pending 微信单主动查单并落态:已付且金额相符→入账(幂等闸),
|
||||
// 渠道关单/超 TTL→过期。返回最新订单 + 是否金额不符(不符则不入账、留人工对账)。
|
||||
// 前端轮询与掉单补偿定时器共用这一份,避免两处「查单→落态」逻辑漂移。
|
||||
func (h *Handler) reconcileOrder(ctx context.Context, o *store.PaymentOrder) (*store.PaymentOrder, bool) {
|
||||
wc := h.pay.Current()
|
||||
if o.Status != store.OrderPending || wc == nil {
|
||||
return o, false
|
||||
}
|
||||
if r, err := wc.QueryOrder(ctx, o.ID); err == nil {
|
||||
switch {
|
||||
case r.Paid && r.AmountFen == o.AmountFen:
|
||||
if _, err := h.db.MarkOrderPaid(ctx, o.ID, r.ChannelTxn); err == nil {
|
||||
o, _ = h.db.GetOrder(ctx, o.ID)
|
||||
}
|
||||
}
|
||||
if o.Status == store.OrderPending && time.Since(o.CreatedAt) > orderTTL {
|
||||
case r.Paid: // 金额对不上:不入账,人工对账(比错账便宜)
|
||||
return o, true
|
||||
case r.Closed:
|
||||
_ = h.db.ExpireOrder(ctx, o.ID)
|
||||
o, _ = h.db.GetOrder(ctx, o.ID)
|
||||
}
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"order": o})
|
||||
if o.Status == store.OrderPending && time.Since(o.CreatedAt) > orderTTL {
|
||||
_ = h.db.ExpireOrder(ctx, o.ID)
|
||||
o, _ = h.db.GetOrder(ctx, o.ID)
|
||||
}
|
||||
return o, false
|
||||
}
|
||||
|
||||
// WechatCallback: POST /api/v1/billing/callback/wechat —— 微信支付回调(公开路由,验签是唯一的门)。
|
||||
@@ -262,3 +276,25 @@ func (h *Handler) AdminPacks(c *gin.Context) {
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"packs": rows})
|
||||
}
|
||||
|
||||
// AdminOrders: GET /api/v1/admin/orders?status= —— 全平台充值订单流 + 状态计数(P5.3 观测)。
|
||||
func (h *Handler) AdminOrders(c *gin.Context) {
|
||||
ctx := c.Request.Context()
|
||||
rows, err := h.db.AllOrders(ctx, c.Query("status"), 50)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"orders": rows, "stats": h.db.OrderStats(ctx)})
|
||||
}
|
||||
|
||||
// AdminReconcile: GET /api/v1/admin/orders/reconcile —— 日终对账(P5.3)。
|
||||
// paid 订单 ↔ 账本 grant 分录逐单比对,列出对不上的(正常应为空)。
|
||||
func (h *Handler) AdminReconcile(c *gin.Context) {
|
||||
rows, err := h.db.ReconcileOrders(c.Request.Context(), 200)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"diffs": rows, "ok": len(rows) == 0})
|
||||
}
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"time"
|
||||
)
|
||||
|
||||
// 掉单补偿(P5.3,设计见 PAYMENT_DESIGN.md §5):
|
||||
// 前端轮询只在「用户开着账单页」时才查单确认——用户扫完码就关页面的话,钱付了、
|
||||
// 订单却永远挂 pending、积分永远不到账。这个后台定时器把「用户在不在场」从入账链路
|
||||
// 里摘掉:周期扫 pending 微信单,逐单 reconcileOrder(与前端轮询同一份幂等落态逻辑)。
|
||||
|
||||
const reconcileInterval = 1 * time.Minute
|
||||
|
||||
// StartReconcile 启动掉单补偿定时器(微信渠道未配置时空转,几乎零成本)。随进程生命周期运行,
|
||||
// ctx 取消即退出。返回给调用方保存以便优雅停机时取消。
|
||||
func (h *Handler) StartReconcile(ctx context.Context) {
|
||||
go func() {
|
||||
t := time.NewTicker(reconcileInterval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-t.C:
|
||||
h.reconcilePending(ctx)
|
||||
}
|
||||
}
|
||||
}()
|
||||
log.Printf("[payment] 掉单补偿定时器已启动(每 %s 扫一次 pending 微信单)", reconcileInterval)
|
||||
}
|
||||
|
||||
// reconcilePending 扫一轮待补偿的 pending 微信单。渠道未配置时直接返回(不打扰)。
|
||||
func (h *Handler) reconcilePending(ctx context.Context) {
|
||||
if h.pay.Current() == nil {
|
||||
return
|
||||
}
|
||||
orders, err := h.db.PendingWechatOrders(ctx, 200)
|
||||
if err != nil {
|
||||
log.Printf("[payment] 补偿扫描取 pending 单失败: %v", err)
|
||||
return
|
||||
}
|
||||
var paid, expired, mismatch int
|
||||
for i := range orders {
|
||||
o := &orders[i]
|
||||
updated, mm := h.reconcileOrder(ctx, o)
|
||||
switch {
|
||||
case mm:
|
||||
mismatch++
|
||||
log.Printf("[payment] ⚠️ 订单 %s 支付金额与订单不符,已挂起待人工对账", o.ID)
|
||||
case updated.Status == "paid":
|
||||
paid++
|
||||
case updated.Status == "expired":
|
||||
expired++
|
||||
}
|
||||
}
|
||||
// 只在有变化时记一行,避免空转刷屏。
|
||||
if paid+expired+mismatch > 0 {
|
||||
log.Printf("[payment] 补偿扫描:入账 %d、过期 %d、金额不符 %d(本轮 %d 单)", paid, expired, mismatch, len(orders))
|
||||
}
|
||||
}
|
||||
@@ -35,6 +35,8 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob.
|
||||
h := handler.New(db, cache, bus, blobStore)
|
||||
// 微信支付渠道装配:DB 配置优先(admin 控制面热重载)→ env 兜底 → 隐藏。失败不阻断启动。
|
||||
h.InitWechat(context.Background())
|
||||
// 掉单补偿定时器:周期扫 pending 微信单确认到账(用户扫完码关页面也能补入账)。
|
||||
h.StartReconcile(context.Background())
|
||||
|
||||
// 可观测性根端点:Prometheus 抓取 + k8s 存活/就绪探针(不挂业务中间件鉴权)。
|
||||
r.GET("/metrics", gin.WrapH(promhttp.Handler()))
|
||||
@@ -140,6 +142,8 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob.
|
||||
admin.PUT("/packs", h.AdminSavePack)
|
||||
admin.GET("/payment/wechat", h.AdminGetWechatPay)
|
||||
admin.PUT("/payment/wechat", h.AdminSaveWechatPay)
|
||||
admin.GET("/orders", h.AdminOrders) // 全平台充值订单流 + 状态计数
|
||||
admin.GET("/orders/reconcile", h.AdminReconcile) // 日终对账:paid 单 ↔ 账本 grant
|
||||
// 多租户成员管理(平台运维口径)
|
||||
admin.GET("/tenants", h.AdminTenants) // 租户目录(成员数+余额)
|
||||
admin.POST("/tenants", h.AdminCreateTenant) // 新建租户(可选指定 owner)
|
||||
|
||||
@@ -249,6 +249,22 @@ func (p *Postgres) MarkOrderPaid(ctx context.Context, orderID, channelTxn string
|
||||
return changed, err
|
||||
}
|
||||
|
||||
// PendingWechatOrders 捞出所有 pending 的微信订单(掉单补偿定时器扫描用)。
|
||||
// 只取微信单:兑换码单核销即 paid,永不 pending,不需要查单。按创建时间升序,先补老单。
|
||||
func (p *Postgres) PendingWechatOrders(ctx context.Context, limit int) ([]PaymentOrder, error) {
|
||||
if p.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
if limit <= 0 || limit > 500 {
|
||||
limit = 200
|
||||
}
|
||||
var out []PaymentOrder
|
||||
err := p.db.WithContext(WithoutTenant(ctx)).
|
||||
Where("status = ? AND channel = ?", OrderPending, ChannelWechat).
|
||||
Order("created_at asc").Limit(limit).Find(&out).Error
|
||||
return out, err
|
||||
}
|
||||
|
||||
// ExpireOrder 把超时未付的 pending 单置为 expired(CAS,已 paid 的不动)。
|
||||
func (p *Postgres) ExpireOrder(ctx context.Context, orderID string) error {
|
||||
if p.db == nil {
|
||||
@@ -259,6 +275,102 @@ func (p *Postgres) ExpireOrder(ctx context.Context, orderID string) error {
|
||||
Update("status", OrderExpired).Error
|
||||
}
|
||||
|
||||
// OrderSummary 是 admin 订单流一行(带租户名,免前端二次查)。
|
||||
type OrderSummary struct {
|
||||
PaymentOrder
|
||||
TenantName string `json:"tenant_name"`
|
||||
}
|
||||
|
||||
// AllOrders 全平台充值订单流(admin 观测;可按状态过滤)。倒序,翻页。
|
||||
func (p *Postgres) AllOrders(ctx context.Context, status string, limit int) ([]OrderSummary, error) {
|
||||
if p.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
if limit <= 0 || limit > 200 {
|
||||
limit = 50
|
||||
}
|
||||
q := p.db.WithContext(WithoutTenant(ctx)).Table("sundynix_payment_order o").
|
||||
Select("o.*, t.name as tenant_name").
|
||||
Joins("LEFT JOIN sundynix_tenant t ON t.id = o.tenant_id").
|
||||
Where("o.deleted_at IS NULL")
|
||||
if status != "" {
|
||||
q = q.Where("o.status = ?", status)
|
||||
}
|
||||
var out []OrderSummary
|
||||
err := q.Order("o.created_at desc").Limit(limit).Scan(&out).Error
|
||||
return out, err
|
||||
}
|
||||
|
||||
// OrderStats 全平台订单状态计数(观测卡片:pending/paid/expired 分布 + 累计到账额)。
|
||||
type OrderStats struct {
|
||||
Pending int64 `json:"pending"`
|
||||
Paid int64 `json:"paid"`
|
||||
Expired int64 `json:"expired"`
|
||||
PaidFenTotal int64 `json:"paid_fen_total"` // 累计到账金额(分),只算真渠道 amount_fen>0
|
||||
}
|
||||
|
||||
func (p *Postgres) OrderStats(ctx context.Context) OrderStats {
|
||||
var s OrderStats
|
||||
if p.db == nil {
|
||||
return s
|
||||
}
|
||||
ctx = WithoutTenant(ctx)
|
||||
// 每次都起新 query builder:复用同一个会累加 WHERE(status=A AND status=B → 恒 0)。
|
||||
countBy := func(status string) int64 {
|
||||
var n int64
|
||||
p.db.WithContext(ctx).Model(&PaymentOrder{}).Where("status = ?", status).Count(&n)
|
||||
return n
|
||||
}
|
||||
s.Pending = countBy(OrderPending)
|
||||
s.Paid = countBy(OrderPaid)
|
||||
s.Expired = countBy(OrderExpired)
|
||||
p.db.WithContext(ctx).Model(&PaymentOrder{}).Where("status = ?", OrderPaid).
|
||||
Select("coalesce(sum(amount_fen),0)").Scan(&s.PaidFenTotal)
|
||||
return s
|
||||
}
|
||||
|
||||
// ReconcileRow 对账差异一行:paid 订单在账本里找不到对应 grant 分录(或反之)。
|
||||
type ReconcileRow struct {
|
||||
OrderID string `json:"order_id"`
|
||||
TenantID string `json:"tenant_id"`
|
||||
CreditsMicro int64 `json:"credits_micro"`
|
||||
Issue string `json:"issue"` // order_without_ledger / ledger_without_order
|
||||
}
|
||||
|
||||
// ReconcileOrders 日终对账:paid 订单 ↔ ledger(kind=grant, ref=订单号) 逐单比对,列出对不上的。
|
||||
// 正常应返回空列表(双闸保证 paid 单必有且仅有一条 grant 分录)。有差异即数据出了问题,需人工查。
|
||||
func (p *Postgres) ReconcileOrders(ctx context.Context, limit int) ([]ReconcileRow, error) {
|
||||
if p.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
if limit <= 0 || limit > 500 {
|
||||
limit = 200
|
||||
}
|
||||
ctx = WithoutTenant(ctx)
|
||||
var out []ReconcileRow
|
||||
// paid 订单但账本无对应 grant 分录(钱记了、积分没到——最严重)
|
||||
if err := p.db.WithContext(ctx).
|
||||
Raw(`SELECT o.id AS order_id, o.tenant_id, o.credits_micro, 'order_without_ledger' AS issue
|
||||
FROM sundynix_payment_order o
|
||||
WHERE o.status='paid' AND o.deleted_at IS NULL
|
||||
AND NOT EXISTS (SELECT 1 FROM sundynix_credit_ledger l
|
||||
WHERE l.kind='grant' AND l.ref=o.id AND l.deleted_at IS NULL)
|
||||
ORDER BY o.created_at DESC LIMIT ?`, limit).Scan(&out).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// grant 分录指向的订单不是 paid(积分到了、订单态不对——重复入账/状态错乱)
|
||||
var out2 []ReconcileRow
|
||||
if err := p.db.WithContext(ctx).
|
||||
Raw(`SELECT l.ref AS order_id, l.tenant_id, l.credits_micro, 'ledger_without_paid_order' AS issue
|
||||
FROM sundynix_credit_ledger l
|
||||
JOIN sundynix_payment_order o ON o.id = l.ref
|
||||
WHERE l.kind='grant' AND l.deleted_at IS NULL AND o.status <> 'paid'
|
||||
ORDER BY l.created_at DESC LIMIT ?`, limit).Scan(&out2).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return append(out, out2...), nil
|
||||
}
|
||||
|
||||
// GetPack 按 id 取在售积分包(下单锁价用;下架的包不可下单)。
|
||||
func (p *Postgres) GetPack(ctx context.Context, id string) (*CreditPack, error) {
|
||||
if p.db == nil {
|
||||
|
||||
Reference in New Issue
Block a user