fix(v2): Runner interval<=0 守卫(skip+WARN)+ 启动预热只跑本地 job(网络型交 ticker 首跳)

- Task 加 Local 字段;interval<=0 时 Add/AddLocal 跳过注册 + WARN 日志(NewTicker(<=0)
  会 panic,Start 的 per-tick recover 覆盖不到该行)。
- 新增 AddLocal/RunOnceLocal:main.go 启动预热改跑 RunOnceLocal,只同步执行纯 DB 任务
  (order-expire/usage-refresh/refund-apply-sweep/refund-stuck-alert);出网 HTTP 的
  sync-pending/paid-spotcheck/crypto-orphan-scan 交各自 ticker 首跳,不再拖住服务启动。
This commit is contained in:
wangjia
2026-07-10 18:12:35 +08:00
parent 9117de7dcf
commit f310c58767
3 changed files with 127 additions and 12 deletions
+41 -4
View File
@@ -10,10 +10,13 @@ import (
"time"
)
// Task 一个周期任务:名字 + 间隔 + 幂等可重跑的 Run。
// Task 一个周期任务:名字 + 间隔 + 幂等可重跑的 Run。Local 标记该任务是否纯本地
// (只碰 DB,无出网 HTTP)——启动预热(RunOnceLocal)只跑 Local 任务,网络型任务交给
// 各自 ticker 首跳,避免上游慢拖住服务启动(main.go r.Run(addr) 之前的同步阶段)。
type Task struct {
Name string
Interval time.Duration
Local bool
Run func(ctx context.Context) error
}
@@ -25,9 +28,29 @@ type Runner struct {
func NewRunner() *Runner { return &Runner{logf: log.Printf} }
// Add 注册一个周期任务
// Add 注册一个周期任务(默认视为网络型/非 Local,启动预热 RunOnceLocal 不跑它,
// 交给 Start 里各自 ticker 的首跳执行)。
func (r *Runner) Add(name string, interval time.Duration, run func(ctx context.Context) error) {
r.tasks = append(r.tasks, Task{Name: name, Interval: interval, Run: run})
r.add(Task{Name: name, Interval: interval, Local: false, Run: run})
}
// AddLocal 注册一个纯本地周期任务(只碰 DB,无出网 HTTP)——会被启动预热
// RunOnceLocal 同步执行一次,让 usage 快照/过期清理/退款自愈等立即生效。
func (r *Runner) AddLocal(name string, interval time.Duration, run func(ctx context.Context) error) {
r.add(Task{Name: name, Interval: interval, Local: true, Run: run})
}
// add 是 Add/AddLocal 的共同落地:Interval<=0 是 operator 配置错误(如
// expire_every_sec: 0)——time.NewTicker 对 <=0 的间隔会 panic,且 Start 里每 tick
// 的 recover 覆盖不到 NewTicker 本身(它在 goroutine 里、ticker 创建那一行就炸,
// 无 defer 保护)。显式优于静默改值:跳过注册 + WARN 日志,而不是偷偷 clamp 成默认值
// 掩盖配置错误。
func (r *Runner) add(t Task) {
if t.Interval <= 0 {
r.logf("[reconcile] WARN 任务 %s interval<=0(%v),跳过注册(检查配置)", t.Name, t.Interval)
return
}
r.tasks = append(r.tasks, t)
}
// exec 跑单个任务一次:panic recover + error 记录,绝不外抛(单任务失败不拖垮其它)。
@@ -42,13 +65,27 @@ func (r *Runner) exec(ctx context.Context, t Task) {
}
}
// RunOnce 顺序跑一遍全部任务(启动预热 + 单测入口)。
// RunOnce 顺序跑一遍全部任务(单测入口;main.go 启动预热改用 RunOnceLocal,
// 网络型任务不再阻塞启动,见下)。
func (r *Runner) RunOnce(ctx context.Context) {
for _, t := range r.tasks {
r.exec(ctx, t)
}
}
// RunOnceLocal 只顺序跑一遍 Local 任务(启动预热用):order-expire/usage-refresh/
// refund-apply-sweep/refund-stuck-alert 等纯 DB 任务立即生效;sync-pending/
// paid-spotcheck/crypto-orphan-scan 等有出网 HTTP(10-15s 超时)的任务跳过,交给
// Start 里各自 ticker 的首跳执行,避免上游慢拖住 r.Run(addr) 前的服务启动。
func (r *Runner) RunOnceLocal(ctx context.Context) {
for _, t := range r.tasks {
if !t.Local {
continue
}
r.exec(ctx, t)
}
}
// Start 每任务一 goroutine + 独立 ticker 常驻;ctx 取消即退出。每 tick 崩溃安全。
func (r *Runner) Start(ctx context.Context) {
for _, t := range r.tasks {
+76
View File
@@ -3,6 +3,7 @@ package reconcile_test
import (
"context"
"errors"
"sync/atomic"
"testing"
"time"
@@ -21,3 +22,78 @@ func TestRunnerRunOnceExecutesAllAndRecoversPanic(t *testing.T) {
t.Fatalf("a=%d b=%d, want 1/1(panic 任务不应阻断其它)", a, b)
}
}
// TestRunnerZeroIntervalSkipsRegistrationAndStartDoesNotPanic 覆盖 Important #1:
// operator 配 interval<=0(如 expire_every_sec: 0)不得让 Start 里 time.NewTicker
// panic 崩进程 —— Add/AddLocal 应在注册期就 skip 该任务(不进 r.tasks),Start 对它
// 不会创建 ticker,自然不 panic;RunOnce/RunOnceLocal 也因任务未注册而不执行它。
func TestRunnerZeroIntervalSkipsRegistrationAndStartDoesNotPanic(t *testing.T) {
r := reconcile.NewRunner()
var zeroRan, negRan, okRan int32
r.Add("zero-interval", 0, func(context.Context) error {
atomic.AddInt32(&zeroRan, 1)
return nil
})
r.AddLocal("neg-interval", -time.Second, func(context.Context) error {
atomic.AddInt32(&negRan, 1)
return nil
})
r.Add("ok-interval", 5*time.Millisecond, func(context.Context) error {
atomic.AddInt32(&okRan, 1)
return nil
})
r.RunOnce(context.Background())
if zeroRan != 0 || negRan != 0 {
t.Fatalf("interval<=0 的任务不应被注册/执行: zeroRan=%d negRan=%d", zeroRan, negRan)
}
if okRan != 1 {
t.Fatalf("正常任务应正常执行一次, got okRan=%d", okRan)
}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
done := make(chan struct{})
go func() {
defer close(done)
r.Start(ctx) // 若 zero-interval/neg-interval 仍被注册,NewTicker(<=0) 会 panic 崩掉这个 goroutine
}()
time.Sleep(30 * time.Millisecond) // 留时间给 ok-interval 的 ticker 至少 tick 一次
cancel()
<-done
if atomic.LoadInt32(&zeroRan) != 0 || atomic.LoadInt32(&negRan) != 0 {
t.Fatalf("Start 之后 interval<=0 的任务仍不应执行: zeroRan=%d negRan=%d", zeroRan, negRan)
}
}
// TestRunnerRunOnceLocalOnlyRunsLocalTasks 覆盖 Important #2:启动预热
// RunOnceLocal 只应执行 AddLocal 注册的任务,Add(默认网络型)注册的任务不跑,
// 交给 Start 里各自 ticker 的首跳。
func TestRunnerRunOnceLocalOnlyRunsLocalTasks(t *testing.T) {
r := reconcile.NewRunner()
var ran []string
r.AddLocal("local-a", time.Minute, func(context.Context) error {
ran = append(ran, "local-a")
return nil
})
r.Add("network-b", time.Minute, func(context.Context) error {
ran = append(ran, "network-b")
return nil
})
r.AddLocal("local-c", time.Minute, func(context.Context) error {
ran = append(ran, "local-c")
return nil
})
r.RunOnceLocal(context.Background())
if len(ran) != 2 {
t.Fatalf("RunOnceLocal 应只跑 2 个 Local 任务, got %v", ran)
}
for _, name := range ran {
if name == "network-b" {
t.Fatalf("RunOnceLocal 不应执行网络型任务, got %v", ran)
}
}
}
+10 -8
View File
@@ -68,23 +68,25 @@ func main() {
rc := config.C.Reconcile
refundStore := store.NewRefundStore(db)
runner := reconcile.NewRunner()
runner.Add("order-expire", time.Duration(rc.ExpireEverySec)*time.Second,
// 纯本地(DB-only)任务用 AddLocal:会被启动预热 RunOnceLocal 同步跑一遍。
runner.AddLocal("order-expire", time.Duration(rc.ExpireEverySec)*time.Second,
reconcile.OrderExpirerTask(orderStore, time.Duration(rc.OrderTTLMin)*time.Minute, time.Now))
runner.Add("usage-refresh", time.Duration(rc.UsageEverySec)*time.Second,
runner.AddLocal("usage-refresh", time.Duration(rc.UsageEverySec)*time.Second,
reconcile.RefreshUsageTask(usage))
runner.AddLocal("refund-apply-sweep", time.Duration(rc.RefundApplyEverySec)*time.Second,
reconcile.RefundApplyTask(orderStore, refundStore, time.Duration(rc.RefundApplyLookbackMin)*time.Minute, time.Now, 200))
runner.AddLocal("refund-stuck-alert", time.Duration(rc.RefundApplyEverySec)*time.Second,
reconcile.RefundStuckAlertTask(refundStore, time.Duration(rc.RefundStuckWarnMin)*time.Minute, time.Now))
// 网络型(出网 HTTP,10-15s 超时)任务用 Add:启动预热不跑,交各自 ticker 首跳。
runner.Add("sync-pending", time.Duration(rc.SyncEverySec)*time.Second,
reconcile.SyncPendingTask(gw, 100))
runner.Add("paid-spotcheck", time.Duration(rc.SpotCheckEverySec)*time.Second,
reconcile.PaidSpotCheckTask(orderStore, pReg, time.Duration(rc.SpotCheckWindowMin)*time.Minute, time.Now))
runner.Add("refund-apply-sweep", time.Duration(rc.RefundApplyEverySec)*time.Second,
reconcile.RefundApplyTask(orderStore, refundStore, time.Duration(rc.RefundApplyLookbackMin)*time.Minute, time.Now, 200))
runner.Add("refund-stuck-alert", time.Duration(rc.RefundApplyEverySec)*time.Second,
reconcile.RefundStuckAlertTask(refundStore, time.Duration(rc.RefundStuckWarnMin)*time.Minute, time.Now))
orphanStore := store.NewOrphanStore(db)
reconcile.AddCryptoJobs(runner, pReg, acctReg, orderStore, orphanStore, rc) // crypto 预留冷启动 Warm + 孤儿扫描(Task 6)
reconcile.AddCryptoJobs(runner, pReg, acctReg, orderStore, orphanStore, rc) // crypto 预留冷启动 Warm(直跑,不受影响)+ 孤儿扫描(网络型,交 ticker 首跳)
ctx := context.Background()
runner.RunOnce(ctx) // 启动预热:先跑一遍(usage 快照/过期清理/退款自愈/crypto Warm 立即生效)
runner.RunOnceLocal(ctx) // 启动预热:只跑本地任务(usage 快照/过期清理/退款自愈立即生效);网络型任务(sync-pending/paid-spotcheck/crypto-orphan-scan)交各自 ticker 首跳,避免上游慢拖住 r.Run(addr) 前的启动
runner.Start(ctx)
log.Printf("[reconcile] 后台守护已启动(过期清理/用量刷新/查单对账/已付抽查/退款修复扫描/退款卡滞告警/crypto 孤儿扫描)")
}