119 lines
2.4 KiB
Go
119 lines
2.4 KiB
Go
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
|
|
}
|