Files
Blizzard f16f63284f feat(cluster): gateway 后台定时器加 leader 锁——多副本安全(B5)
此前两个 gateway 定时器(订阅推进/掉单补偿)每实例各扫一遍:多副本下重复查库,
且对微信查单调用量随副本线性放大(微信有频控,先被限流的是自己)。

store.TryRunExclusive:用 PG advisory try-lock 做集群级单实例执行(leader 选举)——每轮
tick 非阻塞抢锁,抢到才跑、跑完释放;抢不到说明别的实例是 leader、本轮跳过。自愈:锁随
持有连接释放,leader 挂了下一轮别的实例自然抢到接管,无需显式故障转移。订阅/补偿用不同
锁键(可由不同实例分别 lead)。非 PG(sqlite 测试)/无 DB → 退回本地直接跑(单实例安全)。

两个定时器的 tick 各包一层 TryRunExclusive。至此 gateway 可安全多副本(SSE 走共享 Redis
流无需粘性、状态全外置、JWT 无状态,剩数据层 HA 属 C 层 ops)。带 fallback 单测。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 16:21:31 +08:00

112 lines
4.2 KiB
Go

package handler
import (
"context"
"log"
"sync"
"time"
"github.com/sundynix/sundynix-gateway/internal/payment"
)
// 掉单补偿(P5.3,设计见 PAYMENT_DESIGN.md §5):
// 前端轮询只在「用户开着账单页」时才查单确认——用户扫完码就关页面的话,钱付了、
// 订单却永远挂 pending、积分永远不到账。这个后台定时器把「用户在不在场」从入账链路
// 里摘掉:周期扫 pending 微信单,逐单 reconcileOrder(与前端轮询同一份幂等落态逻辑)。
const reconcileInterval = 1 * time.Minute
// reconcileLeaderKey 是掉单补偿的 leader 选举锁键(多副本下只让一个实例查单,见 store.TryRunExclusive)。
const reconcileLeaderKey int64 = 20260723
// 主动查单节流:前端支付弹窗每 2.5s 轮一次单态,而 reconcileOrder 见 pending 就直连渠道
// 查单——单笔订单在 30min TTL 内能打出约 720 次微信查单调用,微信侧有频控,多用户并发时
// 先被限流的反而是我们自己。回调才是入账主路径,查单只是兜底,给它一个最小间隔即可:
// 间隔内的轮询直接返回本地状态(订单一旦被回调入账,本地状态本来就是最新的)。
// 补偿定时器每分钟才跑一轮,远大于这个间隔,不受影响。
const minQueryInterval = 5 * time.Second
// lastQuery: orderID -> 上次真正打渠道查单的时刻。仅用于限流,进程级即可
// (多实例各自限流,量级仍降两个数量级);订单落终态或超期时清理,见 forgetQuery/pruneQueryMarks。
var lastQuery sync.Map
// allowQuery 判断此刻是否放行一次真实查单,放行则记下时刻。
func allowQuery(orderID string) bool {
now := time.Now()
if v, ok := lastQuery.Load(orderID); ok {
if last, _ := v.(time.Time); now.Sub(last) < minQueryInterval {
return false
}
}
lastQuery.Store(orderID, now)
return true
}
func forgetQuery(orderID string) { lastQuery.Delete(orderID) }
// pruneQueryMarks 清掉超过 TTL 的残留标记 —— 用户扫码前就关掉弹窗的订单不会再被轮询,
// 其标记无人清理,不定期回收会随进程运行时长单调增长。
func pruneQueryMarks() {
cutoff := time.Now().Add(-orderTTL)
lastQuery.Range(func(k, v any) bool {
if t, _ := v.(time.Time); t.Before(cutoff) {
lastQuery.Delete(k)
}
return true
})
}
// StartReconcile 启动掉单补偿定时器(微信渠道未配置时空转,几乎零成本)。随进程生命周期运行,
// ctx 取消即退出。返回给调用方保存以便优雅停机时取消。
func (h *Handler) StartReconcile(ctx context.Context) {
safeGo("payment-reconcile-ticker", func() {
t := time.NewTicker(reconcileInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
// 单轮兜底 + leader 选举:多副本下只有抢到锁的实例查单补偿,避免对微信查单量随副本翻倍。
safeCall("payment-reconcile-tick", func() {
h.db.TryRunExclusive(ctx, reconcileLeaderKey, func() {
h.reconcilePending(ctx)
pruneQueryMarks()
})
})
}
}
})
log.Printf("[payment] 掉单补偿定时器已启动(每 %s 扫一次 pending 微信单)", reconcileInterval)
}
// reconcilePending 扫一轮待补偿的 pending 微信单。渠道未配置时直接返回(不打扰)。
func (h *Handler) reconcilePending(ctx context.Context) {
if h.pay.Get(payment.ChannelWechat) == 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))
}
}