package message import ( "context" "log/slog" "time" ) // StartLoops 启动到点分发、到期处理与小时清理循环。推送 worker 在握手时启动。 // 循环随 ctx 取消退出;停机等待由 L-03 调用 WaitLoops。 func (a *App) StartLoops(ctx context.Context) { a.loopWG.Add(3) go func() { defer a.loopWG.Done() a.dispatchLoop(ctx) }() go func() { defer a.loopWG.Done() a.expireLoop(ctx) }() go func() { defer a.loopWG.Done() a.purgeLoop(ctx) }() } // WaitLoops 等待 StartLoops 启动的 goroutine 退出。 func (a *App) WaitLoops() { a.loopWG.Wait() } func (a *App) dispatchLoop(ctx context.Context) { timer := time.NewTimer(time.Millisecond) defer timer.Stop() for { select { case <-ctx.Done(): return case <-a.dispatchCh: case <-timer.C: } opCtx, cancel := context.WithTimeout(ctx, time.Second) nowMs := a.now().UnixMilli() if _, err := a.DispatchDue(opCtx, nowMs, 0); err != nil && ctx.Err() == nil { slog.Error("dispatch due", "err", err) } cancel() delay := time.Second if next, ok := a.earliestScheduled(ctx); ok { d := time.Until(time.UnixMilli(next)) switch { case d < 0: d = 0 case d > time.Minute: d = time.Minute } delay = d } if !timer.Stop() { select { case <-timer.C: default: } } timer.Reset(delay) } } func (a *App) expireLoop(ctx context.Context) { t := time.NewTicker(time.Second) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: opCtx, cancel := context.WithTimeout(ctx, time.Second) if err := a.ExpireOnce(opCtx, a.now().UnixMilli()); err != nil && ctx.Err() == nil { slog.Error("expire once", "err", err) } cancel() } } } func (a *App) purgeLoop(ctx context.Context) { t := time.NewTicker(time.Hour) defer t.Stop() run := func() { opCtx, cancel := context.WithTimeout(ctx, 30*time.Second) if err := a.PurgeOnce(opCtx, a.now().UnixMilli()); err != nil && ctx.Err() == nil { slog.Error("purge once", "err", err) } cancel() } run() for { select { case <-ctx.Done(): return case <-t.C: run() } } } func (a *App) earliestScheduled(ctx context.Context) (int64, bool) { if a.db == nil || a.db.Read == nil { return 0, false } var sendAt int64 err := a.db.Read.QueryRowContext(ctx, ` SELECT MIN(send_at) FROM messages WHERE state = 'scheduled'`).Scan(&sendAt) if err != nil || sendAt == 0 { return 0, false } return sendAt, true }