package handler import ( "context" "log" "time" ) // 订阅推进定时器:周期扫 active 订阅 → 该发的发、该过期的置过期。 // 与掉单补偿(payment_reconcile.go)同一范式:定时器只是"兜底触发器", // 真正的语义与幂等都在 store.TickSubscription 里,两处不会漂移。 // // 为什么需要它:订阅是"有效期内每 N 天发一次积分",没有用户请求来驱动这个节拍。 // 进程停机期间欠下的发放,由 TickSubscription 的补发逻辑一次性补齐。 const subTickInterval = 10 * time.Minute // subLeaderKey 是订阅推进的 leader 选举锁键(多副本下只让一个实例跑,见 store.TryRunExclusive)。 const subLeaderKey int64 = 20260722 // StartSubscriptionTicker 随进程生命周期运行;多实例并发也安全(发放靠 ledger 唯一索引幂等)。 func (h *Handler) StartSubscriptionTicker(ctx context.Context) { safeGo("subscription-ticker", func() { t := time.NewTicker(subTickInterval) defer t.Stop() // 单轮兜底:某轮 panic 不该终止整个定时器,下一轮继续(漏发的下轮补发逻辑兜住)。 // leader 选举:多副本下只有抢到 advisory 锁的实例真正扫,其余跳过(幂等,但省重复扫 + 省频控)。 runTick := func() { safeCall("subscription-tick", func() { h.db.TryRunExclusive(ctx, subLeaderKey, func() { h.tickSubscriptions(ctx) }) }) } runTick() // 启动即跑一次,补停机期间欠的 for { select { case <-ctx.Done(): return case <-t.C: runTick() } } }) log.Printf("[sub] 订阅推进定时器已启动(每 %s 扫一次)", subTickInterval) } func (h *Handler) tickSubscriptions(ctx context.Context) { subs := h.db.DueSubscriptions(ctx, 200) if len(subs) == 0 { return } now := time.Now() var granted, expired int for i := range subs { g, exp, err := h.db.TickSubscription(ctx, &subs[i], now) if err != nil { // 单条失败不影响其它订阅;下一轮会重试(幂等,不会重复发) log.Printf("[sub] ⚠️ 推进订阅 %s 失败: %v", subs[i].ID, err) continue } granted += g if exp { expired++ } } // 只在有变化时记一行,避免空转刷屏 if granted+expired > 0 { log.Printf("[sub] 推进:发放 %d 笔、过期 %d 条(本轮 %d 条订阅)", granted, expired, len(subs)) } }