feat(openflare): 将定时任务迁入 Wavelet Asynq 异步任务框架

将 SSL 续期、可观测数据清理、WAF IP 组同步与 Uptime Kuma 同步
从 API 进程内 cron 迁移为 Asynq Handler,并通过 w_schedules 种子迁移
注册默认定时调度;移除 bootstrap 中的 in-process cron 启动逻辑。
This commit is contained in:
ryan
2026-06-18 19:14:47 +08:00
parent b59a2f3be3
commit 3342d3c39d
15 changed files with 334 additions and 243 deletions
@@ -11,7 +11,6 @@ import (
"time"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/pkg/logger"
)
const (
@@ -29,10 +28,6 @@ var databaseCleanupTargets = map[string]string{
DatabaseCleanupTargetRequestReports: "请求聚合",
}
func init() {
registerJob("database_auto_cleanup", "0 3 * * *", runDatabaseAutoCleanupJob)
}
// DatabaseCleanupInput describes a manual observability cleanup request.
type DatabaseCleanupInput struct {
Target string `json:"target"`
@@ -128,28 +123,6 @@ func RunDatabaseAutoCleanupOnce(now time.Time) (*DatabaseAutoCleanupSummary, err
}, nil
}
func runDatabaseAutoCleanupJob(ctx context.Context) {
summary, err := RunDatabaseAutoCleanupOnce(time.Now())
if err != nil {
logger.ErrorF(ctx, "[OpenFlareTasks] database auto cleanup failed: %v", err)
return
}
if summary == nil {
return
}
totalDeleted := int64(0)
for _, item := range summary.Results {
totalDeleted += item.DeletedCount
}
logger.InfoF(
ctx,
"[OpenFlareTasks] database auto cleanup completed retention_days=%d deleted_count=%d",
summary.RetentionDays,
totalDeleted,
)
}
func deleteAllObservabilityRows(ctx context.Context, target string) (int64, error) {
switch target {
case DatabaseCleanupTargetAccessLogs:
+5 -11
View File
@@ -1,15 +1,9 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
// Package tasks is the single home for OpenFlare scheduled and background work
// that runs inside the API process (goroutines plus robfig/cron), not the Asynq
// worker or scheduler.
// Package tasks hosts shared OpenFlare background job business logic.
//
// Each job lives in its own file and registers via registerJob in init(). This
// layout is intentional so jobs can migrate to a future task framework without
// changing call sites: swap the registry implementation while keeping per-job files.
//
// Wire-up: bootstrap.RegisterOpenFlareBackgroundTasks imports this package so
// init() registrations run; bootstrap.Init starts the cron scheduler when the
// process serves the HTTP API.
package tasks
// Scheduled execution is handled by the Wavelet Asynq task framework; handlers
// live in internal/apps/openflare/async_tasks.go and are registered via
// bootstrap.RegisterTasks().
package tasks
@@ -1,97 +0,0 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package tasks
import (
"context"
"sync"
"github.com/Rain-kl/Wavelet/pkg/logger"
"github.com/robfig/cron/v3"
)
type cronJob struct {
name string
spec string
run func(context.Context)
}
var (
registryMu sync.Mutex
started bool
cronJobs []cronJob
cronRunner *cron.Cron
jobCtx context.Context
jobCancel context.CancelFunc
)
func registerJob(name, cronSpec string, fn func(context.Context)) {
registryMu.Lock()
defer registryMu.Unlock()
cronJobs = append(cronJobs, cronJob{name: name, spec: cronSpec, run: fn})
}
// RegisterCronJob registers a cron job from another OpenFlare package (for example waf).
// Prefer registerJob from init() inside this package when possible.
func RegisterCronJob(name, cronSpec string, fn func(context.Context)) {
registerJob(name, cronSpec, fn)
}
// LogJobError records a failed OpenFlare cron job run.
func LogJobError(ctx context.Context, name string, err error) {
logger.ErrorF(ctx, "[OpenFlareTasks] %s failed: %v", name, err)
}
// Start launches the cron scheduler in a background goroutine. Safe to call multiple times.
func Start(ctx context.Context) {
registryMu.Lock()
defer registryMu.Unlock()
if started {
return
}
jobCtx, jobCancel = context.WithCancel(context.Background())
runner := cron.New()
for _, job := range cronJobs {
current := job
if _, err := runner.AddFunc(current.spec, func() {
current.run(jobCtx)
}); err != nil {
logger.ErrorF(ctx, "[OpenFlareTasks] register cron job %q failed: %v", current.name, err)
continue
}
logger.InfoF(ctx, "[OpenFlareTasks] registered cron job %q (%s)", current.name, current.spec)
}
runner.Start()
cronRunner = runner
started = true
}
// Stop shuts down the cron scheduler gracefully and cancels job contexts.
func Stop() {
registryMu.Lock()
defer registryMu.Unlock()
if !started || cronRunner == nil {
return
}
stopCtx := cronRunner.Stop()
<-stopCtx.Done()
if jobCancel != nil {
jobCancel()
}
cronRunner = nil
started = false
jobCtx = nil
jobCancel = nil
}
// ResetRegistryForTest clears scheduler state so unit tests can call Start again.
func ResetRegistryForTest() {
Stop()
registryMu.Lock()
defer registryMu.Unlock()
cronJobs = nil
}
@@ -12,15 +12,8 @@ import (
"github.com/Rain-kl/Wavelet/pkg/logger"
)
func init() {
registerJob("ssl_renew", "0 0 * * *", func(ctx context.Context) {
if err := runSSLRenewJob(ctx); err != nil {
LogJobError(ctx, "ssl_renew", err)
}
})
}
func runSSLRenewJob(ctx context.Context) error {
// RunSSLRenewJob renews all TLS certificates that are due for renewal.
func RunSSLRenewJob(ctx context.Context) error {
logger.InfoF(ctx, "[OpenFlareTasks] SSL renew job started")
certificates, err := model.ListTLSCertificates(ctx)
@@ -48,4 +41,4 @@ func runSSLRenewJob(ctx context.Context) error {
logger.InfoF(ctx, "[OpenFlareTasks] SSL renew job completed: triggered=%d eligible=%d", triggered, len(due))
return nil
}
}
@@ -64,7 +64,7 @@ func TestRunSSLRenewJobTriggersDueCertificates(t *testing.T) {
require.NoError(t, model.CreateTLSCertificateRecord(ctx, due))
require.NoError(t, model.CreateTLSCertificateRecord(ctx, fresh))
require.NoError(t, runSSLRenewJob(ctx))
require.NoError(t, RunSSLRenewJob(ctx))
renewed, err := model.GetTLSCertificateByID(ctx, due.ID)
require.NoError(t, err)
@@ -1,53 +0,0 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package tasks
import (
"context"
"sync"
"time"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/uptimekuma"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/pkg/logger"
)
var (
lastUptimeKumaSyncTime time.Time
uptimeKumaSyncMutex sync.Mutex
)
func init() {
registerJob("uptime_kuma_sync", "* * * * *", runUptimeKumaSyncJob)
}
func runUptimeKumaSyncJob(ctx context.Context) {
if !model.UptimeKumaEnabled {
return
}
interval := model.UptimeKumaSyncInterval
if interval <= 0 {
interval = 5
}
if time.Since(lastUptimeKumaSyncTime) < time.Duration(interval)*time.Minute {
return
}
if !uptimeKumaSyncMutex.TryLock() {
logger.WarnF(ctx, "[OpenFlareTasks] Uptime Kuma sync job is already running, skipping this scheduled run")
return
}
defer uptimeKumaSyncMutex.Unlock()
logger.InfoF(ctx, "[OpenFlareTasks] Starting scheduled Uptime Kuma sync")
if err := uptimekuma.SyncToUptimeKuma(ctx); err != nil {
logger.ErrorF(ctx, "[OpenFlareTasks] Uptime Kuma sync failed: %v", err)
return
}
lastUptimeKumaSyncTime = time.Now()
logger.InfoF(ctx, "[OpenFlareTasks] Uptime Kuma sync completed successfully")
}
@@ -1,15 +0,0 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package tasks
import "context"
// RegisterWAFIPGroupSync registers the WAF IP group sync cron job without importing waf.
func RegisterWAFIPGroupSync(syncFn func(context.Context) error) {
RegisterCronJob("waf_ip_group_sync", "@every 5m", func(ctx context.Context) {
if err := syncFn(ctx); err != nil {
LogJobError(ctx, "waf_ip_group_sync", err)
}
})
}