mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-11 01:36:37 +08:00
merge: 合并 OpenFlare 定时任务 Asynq 迁移至 dev
This commit is contained in:
@@ -0,0 +1,200 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package openflare
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/apps/openflare/tasks"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/apps/openflare/uptimekuma"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/apps/openflare/waf"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/model"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/task"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
// SSLRenewTask renews due ACME TLS certificates.
|
||||||
|
SSLRenewTask = "openflare:ssl_renew"
|
||||||
|
// TaskTypeSSLRenew is the admin task type for SSL renewal.
|
||||||
|
TaskTypeSSLRenew = "of_ssl_renew"
|
||||||
|
|
||||||
|
// DatabaseAutoCleanupTask prunes observability tables by retention policy.
|
||||||
|
DatabaseAutoCleanupTask = "openflare:database_auto_cleanup"
|
||||||
|
// TaskTypeDatabaseAutoCleanup is the admin task type for observability cleanup.
|
||||||
|
TaskTypeDatabaseAutoCleanup = "of_database_auto_cleanup"
|
||||||
|
|
||||||
|
// WAFIPGroupSyncTask syncs due automatic/subscription WAF IP groups.
|
||||||
|
WAFIPGroupSyncTask = "openflare:waf_ip_group_sync"
|
||||||
|
// TaskTypeWAFIPGroupSync is the admin task type for WAF IP group sync.
|
||||||
|
TaskTypeWAFIPGroupSync = "of_waf_ip_group_sync"
|
||||||
|
|
||||||
|
// UptimeKumaSyncTask synchronizes proxy routes to Uptime Kuma monitors.
|
||||||
|
UptimeKumaSyncTask = "openflare:uptime_kuma_sync"
|
||||||
|
// TaskTypeUptimeKumaSync is the admin task type for Uptime Kuma sync.
|
||||||
|
TaskTypeUptimeKumaSync = "of_uptime_kuma_sync"
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
lastUptimeKumaSyncTime time.Time
|
||||||
|
uptimeKumaSyncMutex sync.Mutex
|
||||||
|
)
|
||||||
|
|
||||||
|
// SSLRenewMeta describes the SSL renewal task.
|
||||||
|
var SSLRenewMeta = task.TaskMeta{
|
||||||
|
Type: TaskTypeSSLRenew,
|
||||||
|
AsynqTask: SSLRenewTask,
|
||||||
|
Name: "OpenFlare SSL 自动续期",
|
||||||
|
Description: "扫描即将到期的 ACME 证书并触发自动续期",
|
||||||
|
SupportsTime: false,
|
||||||
|
MaxRetry: task.DefaultMaxRetry,
|
||||||
|
Queue: task.QueueDefault,
|
||||||
|
Retryable: true,
|
||||||
|
}
|
||||||
|
|
||||||
|
// DatabaseAutoCleanupMeta describes the observability auto-cleanup task.
|
||||||
|
var DatabaseAutoCleanupMeta = task.TaskMeta{
|
||||||
|
Type: TaskTypeDatabaseAutoCleanup,
|
||||||
|
AsynqTask: DatabaseAutoCleanupTask,
|
||||||
|
Name: "OpenFlare 可观测数据自动清理",
|
||||||
|
Description: "按保留天数清理访问日志、性能快照与请求聚合数据",
|
||||||
|
SupportsTime: false,
|
||||||
|
MaxRetry: task.DefaultMaxRetry,
|
||||||
|
Queue: task.QueueDefault,
|
||||||
|
Retryable: true,
|
||||||
|
}
|
||||||
|
|
||||||
|
// WAFIPGroupSyncMeta describes the WAF IP group sync task.
|
||||||
|
var WAFIPGroupSyncMeta = task.TaskMeta{
|
||||||
|
Type: TaskTypeWAFIPGroupSync,
|
||||||
|
AsynqTask: WAFIPGroupSyncTask,
|
||||||
|
Name: "OpenFlare WAF IP 组同步",
|
||||||
|
Description: "同步到期的自动规则与订阅类型 WAF IP 组",
|
||||||
|
SupportsTime: false,
|
||||||
|
MaxRetry: task.DefaultMaxRetry,
|
||||||
|
Queue: task.QueueDefault,
|
||||||
|
Retryable: true,
|
||||||
|
}
|
||||||
|
|
||||||
|
// UptimeKumaSyncMeta describes the Uptime Kuma sync task.
|
||||||
|
var UptimeKumaSyncMeta = task.TaskMeta{
|
||||||
|
Type: TaskTypeUptimeKumaSync,
|
||||||
|
AsynqTask: UptimeKumaSyncTask,
|
||||||
|
Name: "OpenFlare Uptime Kuma 同步",
|
||||||
|
Description: "将启用的代理规则同步到 Uptime Kuma 监控",
|
||||||
|
SupportsTime: false,
|
||||||
|
MaxRetry: task.DefaultMaxRetry,
|
||||||
|
Queue: task.QueueDefault,
|
||||||
|
Retryable: true,
|
||||||
|
}
|
||||||
|
|
||||||
|
// SSLRenewHandler renews due TLS certificates.
|
||||||
|
type SSLRenewHandler struct{}
|
||||||
|
|
||||||
|
// Execute runs SSL certificate renewal for all due certificates.
|
||||||
|
func (h *SSLRenewHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) {
|
||||||
|
task.AppendLog(ctx, "开始扫描待续期证书")
|
||||||
|
if err := tasks.RunSSLRenewJob(ctx); err != nil {
|
||||||
|
task.AppendLog(ctx, "SSL 自动续期失败: %v", err)
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
msg := "SSL 自动续期任务完成"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// DatabaseAutoCleanupHandler prunes observability data when auto-cleanup is enabled.
|
||||||
|
type DatabaseAutoCleanupHandler struct{}
|
||||||
|
|
||||||
|
// Execute runs retention-based cleanup for all observability targets.
|
||||||
|
func (h *DatabaseAutoCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) {
|
||||||
|
if !model.DatabaseAutoCleanupEnabled {
|
||||||
|
msg := "自动清理未启用,跳过执行"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
task.AppendLog(ctx, "开始执行可观测数据自动清理,保留天数=%d", model.DatabaseAutoCleanupRetentionDays)
|
||||||
|
summary, err := tasks.RunDatabaseAutoCleanupOnce(time.Now())
|
||||||
|
if err != nil {
|
||||||
|
task.AppendLog(ctx, "可观测数据自动清理失败: %v", err)
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if summary == nil {
|
||||||
|
msg := "自动清理未启用,跳过执行"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
var totalDeleted int64
|
||||||
|
for _, item := range summary.Results {
|
||||||
|
totalDeleted += item.DeletedCount
|
||||||
|
task.AppendLog(ctx, "清理 %s:删除 %d 条", item.TargetLabel, item.DeletedCount)
|
||||||
|
}
|
||||||
|
|
||||||
|
msg := fmt.Sprintf(
|
||||||
|
"可观测数据自动清理完成,保留 %d 天,共删除 %d 条",
|
||||||
|
summary.RetentionDays,
|
||||||
|
totalDeleted,
|
||||||
|
)
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// WAFIPGroupSyncHandler syncs due WAF IP groups to agents.
|
||||||
|
type WAFIPGroupSyncHandler struct{}
|
||||||
|
|
||||||
|
// Execute syncs all due automatic/subscription WAF IP groups.
|
||||||
|
func (h *WAFIPGroupSyncHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) {
|
||||||
|
task.AppendLog(ctx, "开始同步到期的 WAF IP 组")
|
||||||
|
if err := waf.SyncDueWAFIPGroups(ctx); err != nil {
|
||||||
|
task.AppendLog(ctx, "WAF IP 组同步失败: %v", err)
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
msg := "WAF IP 组同步完成"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// UptimeKumaSyncHandler synchronizes proxy routes to Uptime Kuma.
|
||||||
|
type UptimeKumaSyncHandler struct{}
|
||||||
|
|
||||||
|
// Execute runs Uptime Kuma sync when integration is enabled and the interval has elapsed.
|
||||||
|
func (h *UptimeKumaSyncHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) {
|
||||||
|
if !model.UptimeKumaEnabled {
|
||||||
|
msg := "Uptime Kuma 集成未启用,跳过执行"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
interval := model.UptimeKumaSyncInterval
|
||||||
|
if interval <= 0 {
|
||||||
|
interval = 5
|
||||||
|
}
|
||||||
|
if time.Since(lastUptimeKumaSyncTime) < time.Duration(interval)*time.Minute {
|
||||||
|
msg := fmt.Sprintf("距上次同步不足 %d 分钟,跳过执行", interval)
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
if !uptimeKumaSyncMutex.TryLock() {
|
||||||
|
msg := "Uptime Kuma 同步任务正在执行,跳过本次调度"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
defer uptimeKumaSyncMutex.Unlock()
|
||||||
|
|
||||||
|
task.AppendLog(ctx, "开始同步代理规则到 Uptime Kuma")
|
||||||
|
if err := uptimekuma.SyncToUptimeKuma(ctx); err != nil {
|
||||||
|
task.AppendLog(ctx, "Uptime Kuma 同步失败: %v", err)
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
lastUptimeKumaSyncTime = time.Now()
|
||||||
|
msg := "Uptime Kuma 同步完成"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,83 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package openflare
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/db"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/model"
|
||||||
|
"github.com/glebarez/sqlite"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
"gorm.io/gorm"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestDatabaseAutoCleanupHandlerSkipsWhenDisabled(t *testing.T) {
|
||||||
|
previousEnabled := model.DatabaseAutoCleanupEnabled
|
||||||
|
model.DatabaseAutoCleanupEnabled = false
|
||||||
|
t.Cleanup(func() {
|
||||||
|
model.DatabaseAutoCleanupEnabled = previousEnabled
|
||||||
|
})
|
||||||
|
|
||||||
|
result, err := (&DatabaseAutoCleanupHandler{}).Execute(context.Background(), nil)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NotNil(t, result)
|
||||||
|
assert.Contains(t, result.Message, "未启用")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestDatabaseAutoCleanupHandlerDeletesRowsWhenEnabled(t *testing.T) {
|
||||||
|
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
|
||||||
|
DisableForeignKeyConstraintWhenMigrating: true,
|
||||||
|
})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NoError(t, sqliteDB.AutoMigrate(&model.OpenFlareAccessLog{}))
|
||||||
|
db.SetDB(sqliteDB)
|
||||||
|
t.Cleanup(func() {
|
||||||
|
db.SetDB(nil)
|
||||||
|
})
|
||||||
|
|
||||||
|
now := time.Now().UTC()
|
||||||
|
require.NoError(t, db.DB(context.Background()).Create(&model.OpenFlareAccessLog{
|
||||||
|
NodeID: "node-a",
|
||||||
|
LoggedAt: now.Add(-48 * time.Hour),
|
||||||
|
RemoteAddr: "203.0.113.10",
|
||||||
|
Host: "example.com",
|
||||||
|
Path: "/access",
|
||||||
|
StatusCode: 200,
|
||||||
|
}).Error)
|
||||||
|
|
||||||
|
previousEnabled := model.DatabaseAutoCleanupEnabled
|
||||||
|
previousRetentionDays := model.DatabaseAutoCleanupRetentionDays
|
||||||
|
model.DatabaseAutoCleanupEnabled = true
|
||||||
|
model.DatabaseAutoCleanupRetentionDays = 1
|
||||||
|
t.Cleanup(func() {
|
||||||
|
model.DatabaseAutoCleanupEnabled = previousEnabled
|
||||||
|
model.DatabaseAutoCleanupRetentionDays = previousRetentionDays
|
||||||
|
})
|
||||||
|
|
||||||
|
result, err := (&DatabaseAutoCleanupHandler{}).Execute(context.Background(), nil)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NotNil(t, result)
|
||||||
|
assert.Contains(t, result.Message, "共删除")
|
||||||
|
|
||||||
|
rows, err := model.ListOpenFlareAccessLogs(context.Background(), model.OpenFlareAccessLogQuery{Page: 0, PageSize: 10})
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Empty(t, rows)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestUptimeKumaSyncHandlerSkipsWhenDisabled(t *testing.T) {
|
||||||
|
previousEnabled := model.UptimeKumaEnabled
|
||||||
|
model.UptimeKumaEnabled = false
|
||||||
|
t.Cleanup(func() {
|
||||||
|
model.UptimeKumaEnabled = previousEnabled
|
||||||
|
})
|
||||||
|
|
||||||
|
result, err := (&UptimeKumaSyncHandler{}).Execute(context.Background(), nil)
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NotNil(t, result)
|
||||||
|
assert.Contains(t, result.Message, "未启用")
|
||||||
|
}
|
||||||
@@ -11,7 +11,6 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/model"
|
"github.com/Rain-kl/Wavelet/internal/model"
|
||||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -29,10 +28,6 @@ var databaseCleanupTargets = map[string]string{
|
|||||||
DatabaseCleanupTargetRequestReports: "请求聚合",
|
DatabaseCleanupTargetRequestReports: "请求聚合",
|
||||||
}
|
}
|
||||||
|
|
||||||
func init() {
|
|
||||||
registerJob("database_auto_cleanup", "0 3 * * *", runDatabaseAutoCleanupJob)
|
|
||||||
}
|
|
||||||
|
|
||||||
// DatabaseCleanupInput describes a manual observability cleanup request.
|
// DatabaseCleanupInput describes a manual observability cleanup request.
|
||||||
type DatabaseCleanupInput struct {
|
type DatabaseCleanupInput struct {
|
||||||
Target string `json:"target"`
|
Target string `json:"target"`
|
||||||
@@ -128,28 +123,6 @@ func RunDatabaseAutoCleanupOnce(now time.Time) (*DatabaseAutoCleanupSummary, err
|
|||||||
}, nil
|
}, 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) {
|
func deleteAllObservabilityRows(ctx context.Context, target string) (int64, error) {
|
||||||
switch target {
|
switch target {
|
||||||
case DatabaseCleanupTargetAccessLogs:
|
case DatabaseCleanupTargetAccessLogs:
|
||||||
|
|||||||
@@ -1,15 +1,9 @@
|
|||||||
// Copyright 2026 Arctel.net
|
// Copyright 2026 Arctel.net
|
||||||
// SPDX-License-Identifier: Apache-2.0
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
// Package tasks is the single home for OpenFlare scheduled and background work
|
// Package tasks hosts shared OpenFlare background job business logic.
|
||||||
// that runs inside the API process (goroutines plus robfig/cron), not the Asynq
|
|
||||||
// worker or scheduler.
|
|
||||||
//
|
//
|
||||||
// Each job lives in its own file and registers via registerJob in init(). This
|
// Scheduled execution is handled by the Wavelet Asynq task framework; handlers
|
||||||
// layout is intentional so jobs can migrate to a future task framework without
|
// live in internal/apps/openflare/async_tasks.go and are registered via
|
||||||
// changing call sites: swap the registry implementation while keeping per-job files.
|
// bootstrap.RegisterTasks().
|
||||||
//
|
|
||||||
// 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
|
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"
|
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||||
)
|
)
|
||||||
|
|
||||||
func init() {
|
// RunSSLRenewJob renews all TLS certificates that are due for renewal.
|
||||||
registerJob("ssl_renew", "0 0 * * *", func(ctx context.Context) {
|
func RunSSLRenewJob(ctx context.Context) error {
|
||||||
if err := runSSLRenewJob(ctx); err != nil {
|
|
||||||
LogJobError(ctx, "ssl_renew", err)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
func runSSLRenewJob(ctx context.Context) error {
|
|
||||||
logger.InfoF(ctx, "[OpenFlareTasks] SSL renew job started")
|
logger.InfoF(ctx, "[OpenFlareTasks] SSL renew job started")
|
||||||
|
|
||||||
certificates, err := model.ListTLSCertificates(ctx)
|
certificates, err := model.ListTLSCertificates(ctx)
|
||||||
|
|||||||
@@ -64,7 +64,7 @@ func TestRunSSLRenewJobTriggersDueCertificates(t *testing.T) {
|
|||||||
require.NoError(t, model.CreateTLSCertificateRecord(ctx, due))
|
require.NoError(t, model.CreateTLSCertificateRecord(ctx, due))
|
||||||
require.NoError(t, model.CreateTLSCertificateRecord(ctx, fresh))
|
require.NoError(t, model.CreateTLSCertificateRecord(ctx, fresh))
|
||||||
|
|
||||||
require.NoError(t, runSSLRenewJob(ctx))
|
require.NoError(t, RunSSLRenewJob(ctx))
|
||||||
|
|
||||||
renewed, err := model.GetTLSCertificateByID(ctx, due.ID)
|
renewed, err := model.GetTLSCertificateByID(ctx, due.ID)
|
||||||
require.NoError(t, err)
|
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)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
@@ -1,12 +0,0 @@
|
|||||||
// Copyright 2026 Arctel.net
|
|
||||||
// SPDX-License-Identifier: Apache-2.0
|
|
||||||
|
|
||||||
package waf
|
|
||||||
|
|
||||||
import (
|
|
||||||
oftasks "github.com/Rain-kl/Wavelet/internal/apps/openflare/tasks"
|
|
||||||
)
|
|
||||||
|
|
||||||
func init() {
|
|
||||||
oftasks.RegisterWAFIPGroupSync(SyncDueWAFIPGroups)
|
|
||||||
}
|
|
||||||
@@ -11,8 +11,6 @@ import (
|
|||||||
|
|
||||||
admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push"
|
admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push"
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events"
|
"github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events"
|
||||||
oftasks "github.com/Rain-kl/Wavelet/internal/apps/openflare/tasks"
|
|
||||||
_ "github.com/Rain-kl/Wavelet/internal/apps/openflare/waf"
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/risk_control"
|
"github.com/Rain-kl/Wavelet/internal/apps/risk_control"
|
||||||
taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers"
|
taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers"
|
||||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||||
@@ -25,11 +23,10 @@ type Options struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
registerTasksOnce sync.Once
|
registerTasksOnce sync.Once
|
||||||
registerPushDomainEventsOnce sync.Once
|
registerPushDomainEventsOnce sync.Once
|
||||||
registerTaskListenersOnce sync.Once
|
registerTaskListenersOnce sync.Once
|
||||||
registerOpenFlareBackgroundTasksOnce sync.Once
|
initRuntimeOnce sync.Once
|
||||||
initRuntimeOnce sync.Once
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// RegisterTasks registers all built-in task handlers and metadata.
|
// RegisterTasks registers all built-in task handlers and metadata.
|
||||||
@@ -53,16 +50,10 @@ func RegisterTaskListeners() {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// RegisterOpenFlareBackgroundTasks links OpenFlare in-process cron jobs registered via init().
|
|
||||||
func RegisterOpenFlareBackgroundTasks() {
|
|
||||||
registerOpenFlareBackgroundTasksOnce.Do(func() {})
|
|
||||||
}
|
|
||||||
|
|
||||||
// RegisterAPI wires integrations required by the HTTP API process.
|
// RegisterAPI wires integrations required by the HTTP API process.
|
||||||
func RegisterAPI() {
|
func RegisterAPI() {
|
||||||
RegisterTasks()
|
RegisterTasks()
|
||||||
RegisterPushDomainEvents()
|
RegisterPushDomainEvents()
|
||||||
RegisterOpenFlareBackgroundTasks()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// RegisterWorker wires integrations required by the task worker process.
|
// RegisterWorker wires integrations required by the task worker process.
|
||||||
@@ -81,7 +72,6 @@ func RegisterAll() {
|
|||||||
RegisterTasks()
|
RegisterTasks()
|
||||||
RegisterPushDomainEvents()
|
RegisterPushDomainEvents()
|
||||||
RegisterTaskListeners()
|
RegisterTaskListeners()
|
||||||
RegisterOpenFlareBackgroundTasks()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Init runs shared runtime bootstrap exactly once per process.
|
// Init runs shared runtime bootstrap exactly once per process.
|
||||||
@@ -93,7 +83,6 @@ func Init(ctx context.Context, opts Options) {
|
|||||||
}
|
}
|
||||||
if opts.API {
|
if opts.API {
|
||||||
risk_control.InitLogWriter(ctx)
|
risk_control.InitLogWriter(ctx)
|
||||||
oftasks.Start(ctx)
|
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,11 @@
|
|||||||
|
-- +goose Up
|
||||||
|
INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at)
|
||||||
|
VALUES
|
||||||
|
(101, 'OpenFlare SSL 自动续期', 'of_ssl_renew', '0 0 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||||
|
(102, 'OpenFlare 可观测数据自动清理', 'of_database_auto_cleanup', '0 3 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||||
|
(103, 'OpenFlare WAF IP 组同步', 'of_waf_ip_group_sync', '*/5 * * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||||
|
(104, 'OpenFlare Uptime Kuma 同步', 'of_uptime_kuma_sync', '* * * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||||
|
ON CONFLICT (id) DO NOTHING;
|
||||||
|
|
||||||
|
-- +goose Down
|
||||||
|
DELETE FROM w_schedules WHERE id IN (101, 102, 103, 104);
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
-- +goose Up
|
||||||
|
INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at)
|
||||||
|
VALUES
|
||||||
|
(101, 'OpenFlare SSL 自动续期', 'of_ssl_renew', '0 0 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||||
|
(102, 'OpenFlare 可观测数据自动清理', 'of_database_auto_cleanup', '0 3 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||||
|
(103, 'OpenFlare WAF IP 组同步', 'of_waf_ip_group_sync', '*/5 * * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
|
||||||
|
(104, 'OpenFlare Uptime Kuma 同步', 'of_uptime_kuma_sync', '* * * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
|
||||||
|
ON CONFLICT (id) DO NOTHING;
|
||||||
|
|
||||||
|
-- +goose Down
|
||||||
|
DELETE FROM w_schedules WHERE id IN (101, 102, 103, 104);
|
||||||
@@ -7,6 +7,7 @@ package handlers
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/admin/push"
|
"github.com/Rain-kl/Wavelet/internal/apps/admin/push"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/apps/openflare"
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/upload"
|
"github.com/Rain-kl/Wavelet/internal/apps/upload"
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/user"
|
"github.com/Rain-kl/Wavelet/internal/apps/user"
|
||||||
"github.com/Rain-kl/Wavelet/internal/task"
|
"github.com/Rain-kl/Wavelet/internal/task"
|
||||||
@@ -35,4 +36,17 @@ func Register() {
|
|||||||
// push
|
// push
|
||||||
task.RegisterHandler(push.SendNotificationTask, &push.PushHandler{})
|
task.RegisterHandler(push.SendNotificationTask, &push.PushHandler{})
|
||||||
task.RegisterTaskMeta(push.SendNotificationMeta)
|
task.RegisterTaskMeta(push.SendNotificationMeta)
|
||||||
|
|
||||||
|
// openflare
|
||||||
|
task.RegisterHandler(openflare.SSLRenewTask, &openflare.SSLRenewHandler{})
|
||||||
|
task.RegisterTaskMeta(openflare.SSLRenewMeta)
|
||||||
|
|
||||||
|
task.RegisterHandler(openflare.DatabaseAutoCleanupTask, &openflare.DatabaseAutoCleanupHandler{})
|
||||||
|
task.RegisterTaskMeta(openflare.DatabaseAutoCleanupMeta)
|
||||||
|
|
||||||
|
task.RegisterHandler(openflare.WAFIPGroupSyncTask, &openflare.WAFIPGroupSyncHandler{})
|
||||||
|
task.RegisterTaskMeta(openflare.WAFIPGroupSyncMeta)
|
||||||
|
|
||||||
|
task.RegisterHandler(openflare.UptimeKumaSyncTask, &openflare.UptimeKumaSyncHandler{})
|
||||||
|
task.RegisterTaskMeta(openflare.UptimeKumaSyncMeta)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -40,7 +40,7 @@ sidebar: false
|
|||||||
- 扩展 Relay/Flared heartbeat 载荷与可观测性持久化(frps 观测、健康事件);新增 `of_node_obs_frpc` 单表。
|
- 扩展 Relay/Flared heartbeat 载荷与可观测性持久化(frps 观测、健康事件);新增 `of_node_obs_frpc` 单表。
|
||||||
- Agent heartbeat 恢复 Geo 自动更新、访问日志地域解析与 90 天保留清理;对齐 config `support_files` 过滤规则。
|
- Agent heartbeat 恢复 Geo 自动更新、访问日志地域解析与 90 天保留清理;对齐 config `support_files` 过滤规则。
|
||||||
- 补全 OAuth 快捷路由(`/api/oauth/github`、`/api/oauth/wechat`、`/api/oauth/wechat/bind`、`/api/oauth/email/bind`)。
|
- 补全 OAuth 快捷路由(`/api/oauth/github`、`/api/oauth/wechat`、`/api/oauth/wechat/bind`、`/api/oauth/email/bind`)。
|
||||||
- 新增 `internal/apps/openflare/tasks/` 集中承载 OpenFlare 定时/后台任务(主进程 cron,非 Asynq),含数据库可观测性自动清理、WAF IP 组周期同步、UptimeKuma 同步、ACME 证书自动续期。
|
- 新增 `internal/apps/openflare/tasks/` 集中承载 OpenFlare 定时/后台任务业务逻辑;调度已迁入 Wavelet Asynq 任务框架(`async_tasks.go` + `w_schedules` 种子迁移),含数据库可观测性自动清理、WAF IP 组周期同步、UptimeKuma 同步、ACME 证书自动续期。
|
||||||
- 实装数据库可观测性手动/自动清理、WAF IP 组订阅/自动同步与测试接口、UptimeKuma 监控同步、TLS ACME 申请/续期(lego DNS-01)。
|
- 实装数据库可观测性手动/自动清理、WAF IP 组订阅/自动同步与测试接口、UptimeKuma 监控同步、TLS ACME 申请/续期(lego DNS-01)。
|
||||||
- 修复 Wavelet Agent WebSocket 未处理 `status` 消息导致 WS 模式下 `last_seen_at` 停止更新、节点超时显示离线的问题;列表「最近心跳」恢复显示「WS 已连接」。
|
- 修复 Wavelet Agent WebSocket 未处理 `status` 消息导致 WS 模式下 `last_seen_at` 停止更新、节点超时显示离线的问题;列表「最近心跳」恢复显示「WS 已连接」。
|
||||||
- 修复节点「强制同步」仍为 stub 导致始终返回「节点不在线或通过 WebSocket 发送同步指令失败」的问题。
|
- 修复节点「强制同步」仍为 stub 导致始终返回「节点不在线或通过 WebSocket 发送同步指令失败」的问题。
|
||||||
|
|||||||
Reference in New Issue
Block a user