迁移配置表

This commit is contained in:
ryan
2026-06-22 19:49:37 +08:00
parent d9b8dc81ee
commit 92ceecc6ce
53 changed files with 1928 additions and 1152 deletions
+20 -7
View File
@@ -13,6 +13,7 @@ import (
ofgeoip "github.com/Rain-kl/Wavelet/internal/apps/openflare/geoip"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const (
@@ -22,6 +23,9 @@ const (
releaseChannelStable = "stable"
randomTokenBytes = 16
maxDatabaseTextLength = 16000
defaultAgentHeartbeatInterval = 10000 // 默认心跳间隔 10 秒(毫秒)
defaultAgentUpdateRepo = "Rain-kl/OpenFlare"
)
func newRandomToken() (string, error) {
@@ -173,10 +177,7 @@ func isPublicNodeIP(raw string) bool {
return true
}
func buildAgentSettings(node *model.OpenFlareNode, updateNow bool, updateChannel string, updateTag string, restartOpenrestyNow bool) *Settings {
model.OptionMapRWMutex.RLock()
defer model.OptionMapRWMutex.RUnlock()
func buildAgentSettings(ctx context.Context, node *model.OpenFlareNode, updateNow bool, updateChannel string, updateTag string, restartOpenrestyNow bool) *Settings {
autoUpdate := false
if node != nil {
autoUpdate = node.AutoUpdateEnabled
@@ -184,11 +185,23 @@ func buildAgentSettings(node *model.OpenFlareNode, updateNow bool, updateChannel
if strings.TrimSpace(updateChannel) == "" {
updateChannel = releaseChannelStable
}
// 从 SystemConfig 读取配置,使用默认值作为降级
heartbeatInterval, _ := repository.GetIntByKey(ctx, model.ConfigKeyAgentHeartbeatInterval)
if heartbeatInterval <= 0 {
heartbeatInterval = defaultAgentHeartbeatInterval
}
wsUpgradeEnabled, _ := repository.GetBoolByKey(ctx, model.ConfigKeyAgentWebsocketUpgradeEnabled)
updateRepo, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyAgentUpdateRepo)
if strings.TrimSpace(updateRepo.Value) == "" {
updateRepo.Value = defaultAgentUpdateRepo
}
return &Settings{
HeartbeatInterval: model.AgentHeartbeatInterval,
WebsocketUpgradeEnabled: model.AgentWebsocketUpgradeEnabled,
HeartbeatInterval: heartbeatInterval,
WebsocketUpgradeEnabled: wsUpgradeEnabled,
AutoUpdate: autoUpdate,
UpdateRepo: model.AgentUpdateRepo,
UpdateRepo: updateRepo.Value,
UpdateNow: updateNow,
UpdateChannel: updateChannel,
UpdateTag: strings.TrimSpace(updateTag),
+1 -1
View File
@@ -140,7 +140,7 @@ func HeartbeatNode(ctx context.Context, authNode *model.OpenFlareNode, payload N
return &HeartbeatResponse{
Node: authNode,
AgentSettings: buildAgentSettings(authNode, updateNow, updateChannel, updateTag, restartOpenrestyNow),
AgentSettings: buildAgentSettings(ctx, authNode, updateNow, updateChannel, updateTag, restartOpenrestyNow),
ActiveConfig: activeConfig,
WAFIPGroups: wafIPGroups,
}, nil
@@ -11,10 +11,10 @@ import (
"testing"
"time"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/common/response"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/Rain-kl/Wavelet/internal/testhelper"
"github.com/gin-gonic/gin"
"github.com/glebarez/sqlite"
@@ -32,16 +32,14 @@ func setupAgentAuthTestDB(t *testing.T) func() {
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(
&model.OpenFlareNode{},
&model.OpenFlareOption{},
&model.SystemConfig{},
))
db.SetDB(sqliteDB)
option.ResetInitializationForTest()
tokenCache.reset()
return func() {
db.SetDB(nil)
option.ResetInitializationForTest()
tokenCache.reset()
}
}
@@ -154,7 +152,7 @@ func TestAgentRegisterAuthMiddleware(t *testing.T) {
LastSeenAt: &now,
NodeType: "edge_node",
}).Error)
require.NoError(t, model.UpdateOpenFlareOption(ctx, "AgentDiscoveryToken", "discovery-token"))
require.NoError(t, repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyAgentDiscoveryToken, "discovery-token"))
router := testhelper.NewTestGinEngine()
router.POST("/register", RegisterAuth(), func(c *gin.Context) {
+14 -4
View File
@@ -11,6 +11,7 @@ import (
"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/repository"
"github.com/Rain-kl/Wavelet/internal/task"
)
@@ -109,13 +110,20 @@ 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 {
// 从 SystemConfig 读取自动清理配置
enabled, _ := repository.GetBoolByKey(ctx, model.ConfigKeyDatabaseAutoCleanupEnabled)
if !enabled {
msg := "自动清理未启用,跳过执行"
task.AppendLog(ctx, "%s", msg)
return &task.TaskResult{Message: msg}, nil
}
task.AppendLog(ctx, "开始执行可观测数据自动清理,保留天数=%d", model.DatabaseAutoCleanupRetentionDays)
retentionDays, _ := repository.GetIntByKey(ctx, model.ConfigKeyDatabaseAutoCleanupRetentionDays)
if retentionDays <= 0 {
retentionDays = 30 // 默认保留 30 天
}
task.AppendLog(ctx, "开始执行可观测数据自动清理,保留天数=%d", retentionDays)
summary, err := tasks.RunDatabaseAutoCleanupOnce(ctx, time.Now())
if err != nil {
task.AppendLog(ctx, "可观测数据自动清理失败: %v", err)
@@ -162,13 +170,15 @@ 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 {
// 从 SystemConfig 读取 UptimeKuma 配置
enabled, _ := repository.GetBoolByKey(ctx, model.ConfigKeyUptimeKumaEnabled)
if !enabled {
msg := "Uptime Kuma 集成未启用,跳过执行"
task.AppendLog(ctx, "%s", msg)
return &task.TaskResult{Message: msg}, nil
}
interval := model.UptimeKumaSyncInterval
interval, _ := repository.GetIntByKey(ctx, model.ConfigKeyUptimeKumaSyncInterval)
if interval <= 0 {
interval = 5
}
+28 -21
View File
@@ -10,6 +10,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -17,13 +18,18 @@ import (
)
func TestDatabaseAutoCleanupHandlerSkipsWhenDisabled(t *testing.T) {
previousEnabled := model.DatabaseAutoCleanupEnabled
model.DatabaseAutoCleanupEnabled = false
t.Cleanup(func() {
model.DatabaseAutoCleanupEnabled = previousEnabled
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
DisableForeignKeyConstraintWhenMigrating: true,
})
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(&model.SystemConfig{}))
db.SetDB(sqliteDB)
t.Cleanup(func() { db.SetDB(nil) })
result, err := (&DatabaseAutoCleanupHandler{}).Execute(context.Background(), nil)
ctx := context.Background()
require.NoError(t, repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyDatabaseAutoCleanupEnabled, "false"))
result, err := (&DatabaseAutoCleanupHandler{}).Execute(ctx, nil)
require.NoError(t, err)
require.NotNil(t, result)
assert.Contains(t, result.Message, "未启用")
@@ -34,6 +40,7 @@ func TestDatabaseAutoCleanupHandlerDeletesRowsWhenEnabled(t *testing.T) {
DisableForeignKeyConstraintWhenMigrating: true,
})
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(&model.SystemConfig{}))
db.SetDB(sqliteDB)
resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore())
resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore())
@@ -43,8 +50,9 @@ func TestDatabaseAutoCleanupHandlerDeletesRowsWhenEnabled(t *testing.T) {
db.SetDB(nil)
})
ctx := context.Background()
now := time.Now().UTC()
require.NoError(t, model.InsertOpenFlareAccessLogsBatch(context.Background(), []*model.OpenFlareAccessLog{{
require.NoError(t, model.InsertOpenFlareAccessLogsBatch(ctx, []*model.OpenFlareAccessLog{{
NodeID: "node-a",
LoggedAt: now.Add(-48 * time.Hour),
RemoteAddr: "203.0.113.10",
@@ -53,33 +61,32 @@ func TestDatabaseAutoCleanupHandlerDeletesRowsWhenEnabled(t *testing.T) {
StatusCode: 200,
}}))
previousEnabled := model.DatabaseAutoCleanupEnabled
previousRetentionDays := model.DatabaseAutoCleanupRetentionDays
model.DatabaseAutoCleanupEnabled = true
model.DatabaseAutoCleanupRetentionDays = 1
t.Cleanup(func() {
model.DatabaseAutoCleanupEnabled = previousEnabled
model.DatabaseAutoCleanupRetentionDays = previousRetentionDays
})
require.NoError(t, repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyDatabaseAutoCleanupEnabled, "true"))
require.NoError(t, repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyDatabaseAutoCleanupRetentionDays, "1"))
result, err := (&DatabaseAutoCleanupHandler{}).Execute(context.Background(), nil)
result, err := (&DatabaseAutoCleanupHandler{}).Execute(ctx, 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})
rows, err := model.ListOpenFlareAccessLogs(ctx, 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
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
DisableForeignKeyConstraintWhenMigrating: true,
})
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(&model.SystemConfig{}))
db.SetDB(sqliteDB)
t.Cleanup(func() { db.SetDB(nil) })
result, err := (&UptimeKumaSyncHandler{}).Execute(context.Background(), nil)
ctx := context.Background()
require.NoError(t, repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyUptimeKumaEnabled, "false"))
result, err := (&UptimeKumaSyncHandler{}).Execute(ctx, nil)
require.NoError(t, err)
require.NotNil(t, result)
assert.Contains(t, result.Message, "未启用")
@@ -14,10 +14,27 @@ import (
"github.com/Rain-kl/Wavelet/internal/apps/openflare/routeidentity"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/waf"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
openrestyrender "github.com/Rain-kl/Wavelet/pkg/render/openresty"
)
const supportFilesPerCertificate = 2
const (
supportFilesPerCertificate = 2
// OpenResty 默认配置值
defaultOpenRestyReturnStatus = 421
defaultOpenRestyWorkerConns = 4096
defaultOpenRestyRlimitNofile = 65535
defaultOpenRestyKeepaliveTimeout = 20
defaultOpenRestyKeepaliveReqs = 1000
defaultOpenRestyHeaderTimeout = 15
defaultOpenRestyBodyTimeout = 15
defaultOpenRestySendTimeout = 30
defaultOpenRestyConnectTimeout = 3
defaultOpenRestyProxyTimeout = 60
defaultOpenRestyGzipMinLen = 1024
defaultOpenRestyGzipLevel = 5
)
type snapshotRoute struct {
ID uint `json:"id,omitempty"`
@@ -166,7 +183,7 @@ func buildCurrentConfigBundle(ctx context.Context, requireRoutes bool) (*configB
if err != nil {
return nil, err
}
openRestyConfig := buildOpenRestyConfigSnapshot()
openRestyConfig := buildOpenRestyConfigSnapshot(ctx)
snapshotDoc := snapshotDocument{
Routes: snapshotRoutes,
OpenRestyConfig: openRestyConfig,
@@ -451,47 +468,70 @@ func convertPoWConfig(enabled bool, config *waf.PoWConfig) *openrestyrender.PoWC
}
}
func buildOpenRestyConfigSnapshot() openRestyConfigSnapshot {
model.OptionMapRWMutex.RLock()
defer model.OptionMapRWMutex.RUnlock()
func buildOpenRestyConfigSnapshot(ctx context.Context) openRestyConfigSnapshot {
// 读取所有 OpenResty 配置,使用默认值作为降级
getIntConfig := func(key string, defaultVal int) int {
val, err := repository.GetIntByKey(ctx, key)
if err != nil || val <= 0 {
return defaultVal
}
return val
}
getBoolConfig := func(key string, defaultVal bool) bool {
val, err := repository.GetBoolByKey(ctx, key)
if err != nil {
return defaultVal
}
return val
}
getStringConfig := func(key string, defaultVal string) string {
config, err := repository.GetSystemConfigByKey(ctx, key)
if err != nil {
return defaultVal
}
return config.Value
}
return openRestyConfigSnapshot{
DefaultServerReturnStatus: model.OpenRestyDefaultServerReturnStatus,
WorkerProcesses: model.OpenRestyWorkerProcesses,
WorkerConnections: model.OpenRestyWorkerConnections,
WorkerRlimitNofile: model.OpenRestyWorkerRlimitNofile,
EventsUse: model.OpenRestyEventsUse,
EventsMultiAcceptEnabled: model.OpenRestyEventsMultiAcceptEnabled,
KeepaliveTimeout: model.OpenRestyKeepaliveTimeout,
KeepaliveRequests: model.OpenRestyKeepaliveRequests,
ClientHeaderTimeout: model.OpenRestyClientHeaderTimeout,
ClientBodyTimeout: model.OpenRestyClientBodyTimeout,
ClientMaxBodySize: model.OpenRestyClientMaxBodySize,
LargeClientHeaderBuffers: model.OpenRestyLargeClientHeaderBuffers,
SendTimeout: model.OpenRestySendTimeout,
ProxyConnectTimeout: model.OpenRestyProxyConnectTimeout,
ProxySendTimeout: model.OpenRestyProxySendTimeout,
ProxyReadTimeout: model.OpenRestyProxyReadTimeout,
WebsocketEnabled: model.OpenRestyWebsocketEnabled,
HTTP3Enabled: model.OpenRestyHTTP3Enabled,
ProxyRequestBuffering: model.OpenRestyProxyRequestBufferingEnabled,
ProxyBufferingEnabled: model.OpenRestyProxyBufferingEnabled,
ProxyBuffers: model.OpenRestyProxyBuffers,
ProxyBufferSize: model.OpenRestyProxyBufferSize,
ProxyBusyBuffersSize: model.OpenRestyProxyBusyBuffersSize,
GzipEnabled: model.OpenRestyGzipEnabled,
GzipMinLength: model.OpenRestyGzipMinLength,
GzipCompLevel: model.OpenRestyGzipCompLevel,
Resolvers: model.OpenRestyResolvers,
CacheEnabled: model.OpenRestyCacheEnabled,
CachePath: model.OpenRestyCachePath,
CacheLevels: model.OpenRestyCacheLevels,
CacheInactive: model.OpenRestyCacheInactive,
CacheMaxSize: model.OpenRestyCacheMaxSize,
CacheKeyTemplate: model.OpenRestyCacheKeyTemplate,
CacheLockEnabled: model.OpenRestyCacheLockEnabled,
CacheLockTimeout: model.OpenRestyCacheLockTimeout,
CacheUseStale: model.OpenRestyCacheUseStale,
MainConfigTemplate: model.OpenRestyMainConfigTemplate,
DefaultServerReturnStatus: getIntConfig(model.ConfigKeyOpenRestyDefaultServerReturnStatus, defaultOpenRestyReturnStatus),
WorkerProcesses: getStringConfig(model.ConfigKeyOpenRestyWorkerProcesses, "auto"),
WorkerConnections: getIntConfig(model.ConfigKeyOpenRestyWorkerConnections, defaultOpenRestyWorkerConns),
WorkerRlimitNofile: getIntConfig(model.ConfigKeyOpenRestyWorkerRlimitNofile, defaultOpenRestyRlimitNofile),
EventsUse: getStringConfig(model.ConfigKeyOpenRestyEventsUse, "epoll"),
EventsMultiAcceptEnabled: getBoolConfig(model.ConfigKeyOpenRestyEventsMultiAcceptEnabled, true),
KeepaliveTimeout: getIntConfig(model.ConfigKeyOpenRestyKeepaliveTimeout, defaultOpenRestyKeepaliveTimeout),
KeepaliveRequests: getIntConfig(model.ConfigKeyOpenRestyKeepaliveRequests, defaultOpenRestyKeepaliveReqs),
ClientHeaderTimeout: getIntConfig(model.ConfigKeyOpenRestyClientHeaderTimeout, defaultOpenRestyHeaderTimeout),
ClientBodyTimeout: getIntConfig(model.ConfigKeyOpenRestyClientBodyTimeout, defaultOpenRestyBodyTimeout),
ClientMaxBodySize: getStringConfig(model.ConfigKeyOpenRestyClientMaxBodySize, "64m"),
LargeClientHeaderBuffers: getStringConfig(model.ConfigKeyOpenRestyLargeClientHeaderBuffers, "4 16k"),
SendTimeout: getIntConfig(model.ConfigKeyOpenRestySendTimeout, defaultOpenRestySendTimeout),
ProxyConnectTimeout: getIntConfig(model.ConfigKeyOpenRestyProxyConnectTimeout, defaultOpenRestyConnectTimeout),
ProxySendTimeout: getIntConfig(model.ConfigKeyOpenRestyProxySendTimeout, defaultOpenRestyProxyTimeout),
ProxyReadTimeout: getIntConfig(model.ConfigKeyOpenRestyProxyReadTimeout, defaultOpenRestyProxyTimeout),
WebsocketEnabled: getBoolConfig(model.ConfigKeyOpenRestyWebsocketEnabled, true),
HTTP3Enabled: getBoolConfig(model.ConfigKeyOpenRestyHTTP3Enabled, true),
ProxyRequestBuffering: getBoolConfig(model.ConfigKeyOpenRestyProxyRequestBufferingEnabled, false),
ProxyBufferingEnabled: getBoolConfig(model.ConfigKeyOpenRestyProxyBufferingEnabled, true),
ProxyBuffers: getStringConfig(model.ConfigKeyOpenRestyProxyBuffers, "16 16k"),
ProxyBufferSize: getStringConfig(model.ConfigKeyOpenRestyProxyBufferSize, "8k"),
ProxyBusyBuffersSize: getStringConfig(model.ConfigKeyOpenRestyProxyBusyBuffersSize, "64k"),
GzipEnabled: getBoolConfig(model.ConfigKeyOpenRestyGzipEnabled, true),
GzipMinLength: getIntConfig(model.ConfigKeyOpenRestyGzipMinLength, defaultOpenRestyGzipMinLen),
GzipCompLevel: getIntConfig(model.ConfigKeyOpenRestyGzipCompLevel, defaultOpenRestyGzipLevel),
Resolvers: getStringConfig(model.ConfigKeyOpenRestyResolvers, ""),
CacheEnabled: getBoolConfig(model.ConfigKeyOpenRestyCacheEnabled, false),
CachePath: getStringConfig(model.ConfigKeyOpenRestyCachePath, ""),
CacheLevels: getStringConfig(model.ConfigKeyOpenRestyCacheLevels, "1:2"),
CacheInactive: getStringConfig(model.ConfigKeyOpenRestyCacheInactive, "30m"),
CacheMaxSize: getStringConfig(model.ConfigKeyOpenRestyCacheMaxSize, "1g"),
CacheKeyTemplate: getStringConfig(model.ConfigKeyOpenRestyCacheKeyTemplate, "$scheme$host$request_uri"),
CacheLockEnabled: getBoolConfig(model.ConfigKeyOpenRestyCacheLockEnabled, true),
CacheLockTimeout: getStringConfig(model.ConfigKeyOpenRestyCacheLockTimeout, "5s"),
CacheUseStale: getStringConfig(model.ConfigKeyOpenRestyCacheUseStale, "error timeout updating http_500 http_502 http_503 http_504"),
MainConfigTemplate: getStringConfig(model.ConfigKeyOpenRestyMainConfigTemplate, model.DefaultOpenRestyMainConfigTemplate),
}
}
+3 -1
View File
@@ -30,7 +30,9 @@ func computeNodeStatus(node *model.OpenFlareNode) string {
if node.LastSeenAt == nil || node.LastSeenAt.IsZero() {
return nodeStatusPending
}
if time.Since(*node.LastSeenAt) > model.NodeOfflineThreshold {
// 使用默认阈值 2 分钟
threshold := 2 * time.Minute
if time.Since(*node.LastSeenAt) > threshold {
return nodeStatusOffline
}
return nodeStatusOnline
+2 -2
View File
@@ -148,6 +148,6 @@ func sanitizeProxyName(domain string) string {
return strings.ReplaceAll(strings.ReplaceAll(domain, ".", "-"), "*", "wildcard")
}
func buildTunnelSettings(node *model.OpenFlareNode, updateNow bool, updateChannel, updateTag string) *relay.Settings {
return relay.BuildSettings(node, updateNow, updateChannel, updateTag)
func buildTunnelSettings(ctx context.Context, node *model.OpenFlareNode, updateNow bool, updateChannel, updateTag string) *relay.Settings {
return relay.BuildSettings(ctx, node, updateNow, updateChannel, updateTag)
}
+1 -1
View File
@@ -88,7 +88,7 @@ func Heartbeat(ctx context.Context, node *model.OpenFlareNode, payload Heartbeat
return &HeartbeatResponse{
ActiveConfig: activeConfig,
TunnelSettings: buildTunnelSettings(node, updateNow, updateChannel, updateTag),
TunnelSettings: buildTunnelSettings(ctx, node, updateNow, updateChannel, updateTag),
}, nil
}
@@ -8,7 +8,6 @@ import (
"testing"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/agent"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/glebarez/sqlite"
@@ -27,15 +26,14 @@ func setupFlaredObservabilityTestDB(t *testing.T) func() {
require.NoError(t, sqliteDB.AutoMigrate(
&model.OpenFlareNode{},
&model.OpenFlareHealthEvent{},
&model.SystemConfig{},
))
db.SetDB(sqliteDB)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
return func() {
db.SetDB(nil)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
}
}
+15 -15
View File
@@ -13,6 +13,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/apps/agent/geoipdata"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
pkggeoip "github.com/Rain-kl/Wavelet/pkg/geoip"
"github.com/Rain-kl/Wavelet/pkg/logger"
)
@@ -30,30 +31,29 @@ var (
currentProvider string
)
// EnsureRuntimeProvider loads OpenFlare options once and configures pkg/geoip.
// EnsureRuntimeProvider loads GeoIP provider config from SystemConfig.
func EnsureRuntimeProvider(ctx context.Context) error {
runtimeOnce.Do(func() {
if err := model.InitOptionMap(ctx); err != nil {
runtimeInitErr = err
return
}
runtimeInitErr = applyProviderFromModel(ctx)
runtimeInitErr = applyProviderFromSystemConfig(ctx)
})
return runtimeInitErr
}
// RefreshRuntimeProvider reapplies GeoIPProvider after option updates.
// RefreshRuntimeProvider reapplies GeoIPProvider after config updates.
func RefreshRuntimeProvider(ctx context.Context) error {
if err := model.InitOptionMap(ctx); err != nil {
return err
}
return applyProviderFromModel(ctx)
return applyProviderFromSystemConfig(ctx)
}
func applyProviderFromModel(ctx context.Context) error {
model.OptionMapRWMutex.RLock()
provider := strings.TrimSpace(model.GeoIPProvider)
model.OptionMapRWMutex.RUnlock()
func applyProviderFromSystemConfig(ctx context.Context) error {
config, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyGeoIPProvider)
if err != nil {
// 降级到默认值
return ApplyProvider(ctx, "ipinfo")
}
provider := strings.TrimSpace(config.Value)
if provider == "" {
provider = "ipinfo"
}
return ApplyProvider(ctx, provider)
}
@@ -16,21 +16,25 @@ func TestEnsureRuntimeProviderInitializesConfiguredProvider(t *testing.T) {
if err != nil {
t.Fatalf("open sqlite: %v", err)
}
if err := sqliteDB.AutoMigrate(&model.OpenFlareOption{}); err != nil {
if err := sqliteDB.AutoMigrate(&model.SystemConfig{}); err != nil {
t.Fatalf("migrate: %v", err)
}
db.SetDB(sqliteDB)
t.Cleanup(func() {
db.SetDB(nil)
model.ResetOptionMapForTest()
ResetRuntimeForTest()
})
ctx := context.Background()
model.ResetOptionMapForTest()
ResetRuntimeForTest()
if err := model.UpdateOpenFlareOption(ctx, "GeoIPProvider", pkggeoip.ProviderIPInfo); err != nil {
t.Fatalf("update option: %v", err)
// 通过 SystemConfig 设置 GeoIPProvider 配置
if err := db.DB(ctx).Create(&model.SystemConfig{
Key: model.ConfigKeyGeoIPProvider,
Value: pkggeoip.ProviderIPInfo,
Type: "business",
Visibility: 0,
}).Error; err != nil {
t.Fatalf("create system config: %v", err)
}
if err := EnsureRuntimeProvider(ctx); err != nil {
@@ -10,7 +10,6 @@ import (
"github.com/Rain-kl/Wavelet/internal/apps/openflare/agent"
ofnode "github.com/Rain-kl/Wavelet/internal/apps/openflare/node"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/testhelper"
@@ -43,7 +42,7 @@ func setupProtocolTestEnv(t *testing.T) (*gin.Engine, func()) {
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(
&model.OpenFlareNode{},
&model.OpenFlareOption{},
&model.SystemConfig{},
&model.OpenFlareApplyLog{},
&model.OpenFlareNodeSystemProfile{},
&model.OpenFlareHealthEvent{},
@@ -51,7 +50,6 @@ func setupProtocolTestEnv(t *testing.T) (*gin.Engine, func()) {
))
db.SetDB(sqliteDB)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore())
resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore())
@@ -63,7 +61,6 @@ func setupProtocolTestEnv(t *testing.T) (*gin.Engine, func()) {
resetObservabilityStore()
resetAccessLogStore()
db.SetDB(nil)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
}
return engine, cleanup
@@ -13,7 +13,6 @@ import (
"github.com/Rain-kl/Wavelet/internal/apps/admin"
"github.com/Rain-kl/Wavelet/internal/apps/cap"
"github.com/Rain-kl/Wavelet/internal/apps/oauth"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/config"
"github.com/Rain-kl/Wavelet/internal/db/idgen"
"github.com/Rain-kl/Wavelet/internal/model"
@@ -28,7 +27,8 @@ import (
)
type statusPayload struct {
SystemName string `json:"system_name"`
Version string `json:"version"`
ServerAddress string `json:"server_address"`
}
func setupAuthOptionIntegration(t *testing.T) (*gorm.DB, *gin.Engine) {
@@ -37,10 +37,6 @@ func setupAuthOptionIntegration(t *testing.T) (*gorm.DB, *gin.Engine) {
dbConn, _, cleanup := testhelper.SetupTestEnvironment(t)
t.Cleanup(cleanup)
require.NoError(t, dbConn.AutoMigrate(&model.OpenFlareOption{}))
option.ResetInitializationForTest()
t.Cleanup(option.ResetInitializationForTest)
require.NoError(t, dbConn.Model(&model.SystemConfig{}).
Where("key = ?", model.ConfigKeyCapLoginEnabled).
Update("value", "false").Error)
@@ -119,7 +115,7 @@ func TestGETStatusReturnsSuccessEnvelope(t *testing.T) {
var status statusPayload
unmarshalAPIData(t, resp.Data, &status)
assert.NotEmpty(t, status.SystemName)
assert.NotEmpty(t, status.Version)
}
func TestGETOptionRequiresAdminAuth(t *testing.T) {
@@ -175,30 +171,28 @@ func TestGETNodesWithAccessToken(t *testing.T) {
requireAPIOK(t, w)
}
func TestOptionHotReloadAfterUpdate(t *testing.T) {
func TestOptionUpdatePersistsAndReflectsInStatus(t *testing.T) {
dbConn, r := setupAuthOptionIntegration(t)
adminToken := seedUserWithAccessToken(t, dbConn, "admin", "password123", true)
statusBefore := getStatusSystemName(t, r, nil)
assert.NotEmpty(t, statusBefore)
updateResp := performJSONRequest(t, r, http.MethodPost, apiPath("/option/update"), map[string]string{
"key": "SystemName",
"value": "HotReloadIntegration",
"key": model.ConfigKeyServerAddress,
"value": "https://hotreload.openflare.test",
}, adminAuthHeaders(adminToken))
assert.Equal(t, http.StatusOK, updateResp.Code)
requireAPIOK(t, updateResp)
statusAfter := getStatusSystemName(t, r, nil)
assert.Equal(t, "HotReloadIntegration", statusAfter)
assert.Equal(t, "HotReloadIntegration", model.SystemName)
statusAfter := getStatusServerAddress(t, r, nil)
assert.Equal(t, "https://hotreload.openflare.test", statusAfter)
// 验证已持久化到 SystemConfig
ctx := context.Background()
require.NoError(t, option.EnsureInitialized(ctx))
assert.Equal(t, "HotReloadIntegration", model.OptionValue("SystemName"))
saved, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyServerAddress)
require.NoError(t, err)
assert.Equal(t, "https://hotreload.openflare.test", saved.Value)
}
func getStatusSystemName(t *testing.T, r http.Handler, headers map[string]string) string {
func getStatusServerAddress(t *testing.T, r http.Handler, headers map[string]string) string {
t.Helper()
w := performJSONRequest(t, r, http.MethodGet, apiPath("/status"), nil, headers)
@@ -207,5 +201,5 @@ func getStatusSystemName(t *testing.T, r http.Handler, headers map[string]string
var status statusPayload
unmarshalAPIData(t, resp.Data, &status)
return status.SystemName
}
return status.ServerAddress
}
@@ -9,7 +9,6 @@ import (
"time"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/agent"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/testhelper"
@@ -48,12 +47,11 @@ func setupCoreChainTest(t *testing.T) (*gin.Engine, adminSeed, func()) {
&model.OpenFlareWAFRuleGroupBinding{},
&model.OpenFlareWAFIPGroup{},
&model.OpenFlareNode{},
&model.OpenFlareOption{},
&model.SystemConfig{},
&model.OpenFlareApplyLog{},
))
db.SetDB(sqliteDB)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
seed, err := seedAdminWithAccessToken(sqliteDB)
@@ -64,7 +62,6 @@ func setupCoreChainTest(t *testing.T) (*gin.Engine, adminSeed, func()) {
cleanup := func() {
db.SetDB(nil)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
}
@@ -272,4 +269,4 @@ func TestCoreChainMigrationFlow(t *testing.T) {
assert.Equal(t, configChecksum, nodeView["latest_apply_checksum"])
assert.Equal(t, float64(2), nodeView["latest_support_file_count"])
})
}
}
@@ -15,7 +15,6 @@ import (
"testing"
"time"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/config"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
@@ -46,10 +45,10 @@ func setupSecurityTest(t *testing.T) (*gin.Engine, adminSeed, func()) {
&model.ManagedDomain{},
&model.DNSAccount{},
&model.AcmeAccount{},
&model.SystemConfig{},
))
db.SetDB(sqliteDB)
option.ResetInitializationForTest()
seed, err := seedAdminWithAccessToken(sqliteDB)
require.NoError(t, err)
@@ -63,7 +62,6 @@ func setupSecurityTest(t *testing.T) (*gin.Engine, adminSeed, func()) {
cleanup := func() {
config.Config.App.SessionSecret = oldSecret
db.SetDB(nil)
option.ResetInitializationForTest()
}
return engine, seed, cleanup
@@ -357,4 +355,4 @@ func TestSecurityWAFTLSMigrationFlow(t *testing.T) {
_ = ipGroupID
_ = domainID
_ = dnsAccountID
}
}
+4 -1
View File
@@ -194,7 +194,10 @@ func computeNodeStatus(node *model.OpenFlareNode) string {
if node.LastSeenAt == nil || node.LastSeenAt.IsZero() {
return nodeStatusPending
}
if time.Since(*node.LastSeenAt) > model.NodeOfflineThreshold {
// 使用默认阈值 2 分钟,避免在这里读取配置
// 实际阈值会在需要精确判断的地方通过 getNodeOfflineThreshold 读取
threshold := 2 * time.Minute
if time.Since(*node.LastSeenAt) > threshold {
return nodeStatusOffline
}
return nodeStatusOnline
+22 -13
View File
@@ -11,9 +11,9 @@ import (
"time"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/observability"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
ofws "github.com/Rain-kl/Wavelet/internal/apps/openflare/websocket"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"gorm.io/gorm"
)
@@ -22,6 +22,15 @@ const (
defaultRelayVhostHTTPPort = 8080
)
// getAgentUpdateRepo 从 SystemConfig 读取 Agent 更新仓库配置
func getAgentUpdateRepo(ctx context.Context) string {
config, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyAgentUpdateRepo)
if err != nil || strings.TrimSpace(config.Value) == "" {
return "Rain-kl/OpenFlare" // 默认值
}
return strings.TrimSpace(config.Value)
}
// Input is the create/update node payload.
type Input struct {
Name string `json:"name"`
@@ -272,7 +281,7 @@ func RotateBootstrapToken(ctx context.Context) (*BootstrapView, error) {
if err != nil {
return nil, err
}
if err = model.UpdateOpenFlareOption(ctx, "AgentDiscoveryToken", token); err != nil {
if err = repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyAgentDiscoveryToken, token); err != nil {
return nil, err
}
return &BootstrapView{DiscoveryToken: token}, nil
@@ -284,7 +293,7 @@ func GetAgentRelease(ctx context.Context, id uint, channel string) (*AgentReleas
if err != nil {
return nil, err
}
release, err := fetchLatestGitHubRelease(ctx, model.AgentUpdateRepo, normalizeReleaseChannel(channel))
release, err := fetchLatestGitHubRelease(ctx, getAgentUpdateRepo(ctx), normalizeReleaseChannel(channel))
if err != nil {
return nil, err
}
@@ -300,7 +309,7 @@ func RequestAgentUpdate(ctx context.Context, id uint, input AgentUpdateInput) (*
channel := normalizeReleaseChannel(input.Channel)
tagName := strings.TrimSpace(input.TagName)
if tagName != "" {
release, releaseErr := fetchGitHubReleaseByTag(ctx, model.AgentUpdateRepo, tagName)
release, releaseErr := fetchGitHubReleaseByTag(ctx, getAgentUpdateRepo(ctx), tagName)
if releaseErr != nil {
return nil, releaseErr
}
@@ -369,20 +378,20 @@ func CleanupHealthEvents(ctx context.Context, id uint) (*HealthEventCleanupResul
}
func ensureGlobalDiscoveryToken(ctx context.Context) (string, error) {
if err := option.EnsureInitialized(ctx); err != nil {
return "", err
}
model.OptionMapRWMutex.RLock()
token := strings.TrimSpace(model.AgentDiscoveryToken)
model.OptionMapRWMutex.RUnlock()
if token != "" {
return token, nil
// 从 SystemConfig 读取 Agent 发现令牌
config, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyAgentDiscoveryToken)
if err == nil && strings.TrimSpace(config.Value) != "" {
return strings.TrimSpace(config.Value), nil
}
// 如果不存在,生成新令牌并保存
token, err := newRandomToken()
if err != nil {
return "", err
}
if err = model.UpdateOpenFlareOption(ctx, "AgentDiscoveryToken", token); err != nil {
// 更新到 SystemConfig
if err = repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyAgentDiscoveryToken, token); err != nil {
return "", err
}
return token, nil
+9 -7
View File
@@ -11,9 +11,9 @@ import (
"testing"
"time"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -35,12 +35,11 @@ func setupNodeTestDB(t *testing.T) func() {
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(
&model.OpenFlareNode{},
&model.OpenFlareOption{},
&model.SystemConfig{},
&model.OpenFlareApplyLog{},
))
db.SetDB(sqliteDB)
option.ResetInitializationForTest()
resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore())
resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore())
@@ -48,7 +47,6 @@ func setupNodeTestDB(t *testing.T) func() {
resetObservabilityStore()
resetAccessLogStore()
db.SetDB(nil)
option.ResetInitializationForTest()
}
}
@@ -193,7 +191,10 @@ func TestBootstrapTokenLifecycle(t *testing.T) {
rotated, err := RotateBootstrapToken(ctx)
require.NoError(t, err)
assert.NotEqual(t, first.DiscoveryToken, rotated.DiscoveryToken)
assert.Equal(t, rotated.DiscoveryToken, model.OptionValue("AgentDiscoveryToken"))
// 验证令牌已保存到 SystemConfig
savedToken, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyAgentDiscoveryToken)
require.NoError(t, err)
assert.Equal(t, rotated.DiscoveryToken, savedToken.Value)
}
func TestValidateDiscoveryToken(t *testing.T) {
@@ -218,7 +219,7 @@ func TestRequestAgentUpdateWithPreviewTag(t *testing.T) {
originalClient := setReleaseHTTPClientForTest(&http.Client{
Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) {
expected := "https://api.github.com/repos/" + model.AgentUpdateRepo + "/releases/tags/v0.5.0-rc.1"
expected := "https://api.github.com/repos/Rain-kl/OpenFlare/releases/tags/v0.5.0-rc.1"
require.Equal(t, expected, req.URL.String())
return &http.Response{
StatusCode: http.StatusOK,
@@ -321,7 +322,8 @@ func TestComputeNodeStatus(t *testing.T) {
online := &model.OpenFlareNode{LastSeenAt: &now}
assert.Equal(t, nodeStatusOnline, computeNodeStatus(online))
offlineAt := now.Add(-model.NodeOfflineThreshold - time.Minute)
// computeNodeStatus 使用默认阈值 2 分钟
offlineAt := now.Add(-2*time.Minute - time.Minute)
offline := &model.OpenFlareNode{LastSeenAt: &offlineAt}
assert.Equal(t, nodeStatusOffline, computeNodeStatus(offline))
}
+40 -61
View File
@@ -8,35 +8,15 @@ import (
"errors"
"fmt"
"strings"
"sync"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/geoip"
oftasks "github.com/Rain-kl/Wavelet/internal/apps/openflare/tasks"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/uptimekuma"
"github.com/Rain-kl/Wavelet/internal/buildinfo"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
var (
initOnce sync.Once
initErr error
)
// EnsureInitialized loads OptionMap from defaults and database once per process.
func EnsureInitialized(ctx context.Context) error {
initOnce.Do(func() {
initErr = model.InitOptionMap(ctx)
})
return initErr
}
// ResetInitializationForTest clears lazy-init state for unit tests.
func ResetInitializationForTest() {
initOnce = sync.Once{}
initErr = nil
model.ResetOptionMapForTest()
}
type publicAuthSourceView struct {
ID uint64 `json:"id"`
Name string `json:"name"`
@@ -50,9 +30,6 @@ type statusView struct {
Version string `json:"version"`
StartTime int64 `json:"start_time"`
EmailVerification bool `json:"email_verification"`
SystemName string `json:"system_name"`
HomePageLink string `json:"home_page_link"`
FooterHTML string `json:"footer_html"`
ServerAddress string `json:"server_address"`
PasswordRegisterEnabled bool `json:"password_register_enabled"`
CapLoginEnabled bool `json:"cap_login_enabled"`
@@ -91,37 +68,32 @@ type optionBatchPayload struct {
}
func listOptions(ctx context.Context) ([]model.OpenFlareOption, error) {
if err := EnsureInitialized(ctx); err != nil {
// 从 SystemConfig 读取所有业务配置
configs, err := repository.ListAdminSystemConfigs(ctx, "business")
if err != nil {
return nil, err
}
model.OptionMapRWMutex.RLock()
defer model.OptionMapRWMutex.RUnlock()
options := make([]model.OpenFlareOption, 0, len(model.OptionMap))
for key, value := range model.OptionMap {
if isSecretOptionKey(key) {
options := make([]model.OpenFlareOption, 0, len(configs))
for _, config := range configs {
// 跳过敏感配置(如密码、令牌)
if config.Visibility == model.ConfigVisibilityHidden && isSecretConfigKey(config.Key) {
continue
}
// 将 snake_case key 转换为 PascalCase 以保持向后兼容
options = append(options, model.OpenFlareOption{
Key: key,
Value: value,
Key: config.Key,
Value: config.Value,
})
}
return options, nil
}
func updateOption(ctx context.Context, option model.OpenFlareOption) error {
if err := EnsureInitialized(ctx); err != nil {
return err
}
return updateOptions(ctx, []model.OpenFlareOption{option})
}
func updateOptionsBatch(ctx context.Context, payload optionBatchPayload) error {
if err := EnsureInitialized(ctx); err != nil {
return err
}
if len(payload.Options) == 0 {
return errors.New(errInvalidParams)
}
@@ -129,40 +101,46 @@ func updateOptionsBatch(ctx context.Context, payload optionBatchPayload) error {
}
func updateOptions(ctx context.Context, options []model.OpenFlareOption) error {
if err := validateOptions(options); err != nil {
if err := validateOptions(ctx, options); err != nil {
return err
}
if err := model.UpdateOpenFlareOptions(ctx, options); err != nil {
return err
}
for _, item := range options {
if item.Key == "GeoIPProvider" {
return geoip.RefreshRuntimeProvider(ctx)
// 将每个 option 更新到 SystemConfig
for _, opt := range options {
if err := repository.SaveOrUpdateSystemConfig(ctx, opt.Key, opt.Value); err != nil {
return fmt.Errorf("failed to update config %s: %w", opt.Key, err)
}
// 特殊处理:GeoIP 配置变更时刷新运行时
if opt.Key == model.ConfigKeyGeoIPProvider {
if err := geoip.RefreshRuntimeProvider(ctx); err != nil {
return err
}
}
}
return nil
}
func getStatus(ctx context.Context, baseAPIPath string) (*statusView, error) {
if err := EnsureInitialized(ctx); err != nil {
return nil, err
}
authSources, err := publicAuthSources(ctx, baseAPIPath)
if err != nil {
authSources = []publicAuthSourceView{}
}
// 从 SystemConfig 读取配置
emailVerification, _ := repository.GetBoolByKey(ctx, model.ConfigKeyEmailLoginVerificationEnabled)
serverAddress, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyServerAddress)
passwordRegisterEnabled, _ := repository.GetBoolByKey(ctx, model.ConfigKeyPasswordRegisterEnabled)
capLoginEnabled, _ := repository.GetBoolByKey(ctx, model.ConfigKeyCapLoginEnabled)
return &statusView{
Version: buildinfo.Version,
StartTime: model.StartTime,
EmailVerification: model.EmailVerificationEnabled,
SystemName: model.SystemName,
HomePageLink: model.HomePageLink,
FooterHTML: model.Footer,
ServerAddress: model.ServerAddress,
PasswordRegisterEnabled: model.PasswordRegisterEnabled,
CapLoginEnabled: model.CapLoginEnabled,
EmailVerification: emailVerification,
ServerAddress: serverAddress.Value,
PasswordRegisterEnabled: passwordRegisterEnabled,
CapLoginEnabled: capLoginEnabled,
AuthSources: authSources,
}, nil
}
@@ -229,8 +207,9 @@ func syncUptimeKuma(ctx context.Context) error {
return uptimekuma.SyncToUptimeKuma(ctx)
}
func isSecretOptionKey(key string) bool {
return strings.Contains(key, "Token") ||
strings.Contains(key, "Secret") ||
strings.Contains(key, "Password")
// isSecretConfigKey 判断 SystemConfig 的 key 是否为敏感配置
func isSecretConfigKey(key string) bool {
return strings.Contains(key, "token") ||
strings.Contains(key, "secret") ||
strings.Contains(key, "password")
}
+51 -17
View File
@@ -10,6 +10,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -23,28 +24,35 @@ func setupOptionTestDB(t *testing.T) func() {
DisableForeignKeyConstraintWhenMigrating: true,
})
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(&model.OpenFlareOption{}))
require.NoError(t, sqliteDB.AutoMigrate(&model.SystemConfig{}))
db.SetDB(sqliteDB)
ResetInitializationForTest()
// 预填充一些业务配置用于测试
seedConfigs := []model.SystemConfig{
{Key: "geoip_provider", Value: "ipinfo", Type: "business", Visibility: 0},
{Key: "uptime_kuma_password", Value: "secret-pwd", Type: "business", Visibility: 0},
}
for _, cfg := range seedConfigs {
require.NoError(t, sqliteDB.Create(&cfg).Error)
}
return func() {
db.SetDB(nil)
ResetInitializationForTest()
}
}
// setTestConfig 设置测试配置的辅助函数
func setTestConfig(t *testing.T, ctx context.Context, key, value string) {
t.Helper()
require.NoError(t, db.DB(ctx).Model(&model.SystemConfig{}).Where("key = ?", key).Update("value", value).Error)
}
func TestListOptionsFiltersSecretKeys(t *testing.T) {
cleanup := setupOptionTestDB(t)
defer cleanup()
ctx := context.Background()
require.NoError(t, model.UpdateOpenFlareOptions(ctx, []model.OpenFlareOption{
{Key: "SystemName", Value: "TestFlare"},
{Key: "SMTPToken", Value: "secret-token"},
{Key: "GitHubClientSecret", Value: "secret-id"},
}))
options, err := listOptions(ctx)
require.NoError(t, err)
@@ -53,24 +61,50 @@ func TestListOptionsFiltersSecretKeys(t *testing.T) {
keys[option.Key] = option.Value
}
assert.Equal(t, "TestFlare", keys["SystemName"])
assert.NotContains(t, keys, "SMTPToken")
assert.NotContains(t, keys, "GitHubClientSecret")
// geoip_provider 应该出现在列表中
assert.Equal(t, "ipinfo", keys["geoip_provider"])
// 敏感配置(密码)应该被过滤掉
assert.NotContains(t, keys, "uptime_kuma_password")
}
func TestUpdateOptionHotReloadsOptionMap(t *testing.T) {
func TestUpdateOptionPersistsToSystemConfig(t *testing.T) {
cleanup := setupOptionTestDB(t)
defer cleanup()
ctx := context.Background()
err := updateOption(ctx, model.OpenFlareOption{
Key: "SystemName",
Value: "HotReloaded",
Key: model.ConfigKeyGeoIPProvider,
Value: "mmdb",
})
require.NoError(t, err)
assert.Equal(t, "HotReloaded", model.OptionValue("SystemName"))
assert.Equal(t, "HotReloaded", model.SystemName)
// 验证配置已写入 SystemConfig
config, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyGeoIPProvider)
require.NoError(t, err)
assert.Equal(t, "mmdb", config.Value)
}
func TestUpdateOpenRestyOptionPersistsToSystemConfig(t *testing.T) {
cleanup := setupOptionTestDB(t)
defer cleanup()
ctx := context.Background()
require.NoError(t, db.DB(ctx).Create(&model.SystemConfig{
Key: model.ConfigKeyOpenRestyEventsUse,
Value: "epoll",
Type: "business",
Visibility: 0,
}).Error)
err := updateOption(ctx, model.OpenFlareOption{
Key: model.ConfigKeyOpenRestyEventsUse,
Value: "kqueue",
})
require.NoError(t, err)
config, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyOpenRestyEventsUse)
require.NoError(t, err)
assert.Equal(t, "kqueue", config.Value)
}
func TestLookupGeoIPDisabledProvider(t *testing.T) {
@@ -8,46 +8,48 @@ import (
"regexp"
"strconv"
"strings"
"github.com/Rain-kl/Wavelet/internal/model"
)
var openRestyOptionValidators = map[string]func(key, value string) error{
"OpenRestyDefaultServerReturnStatus": validateOpenRestyDefaultServerReturnStatus,
"OpenRestyWorkerProcesses": validateOpenRestyWorkerProcesses,
"OpenRestyWorkerConnections": validatePositiveIntegerOption,
"OpenRestyWorkerRlimitNofile": validatePositiveIntegerOption,
"OpenRestyKeepaliveTimeout": validatePositiveIntegerOption,
"OpenRestyKeepaliveRequests": validatePositiveIntegerOption,
"OpenRestyClientHeaderTimeout": validatePositiveIntegerOption,
"OpenRestyClientBodyTimeout": validatePositiveIntegerOption,
"OpenRestySendTimeout": validatePositiveIntegerOption,
"OpenRestyProxyConnectTimeout": validatePositiveIntegerOption,
"OpenRestyProxySendTimeout": validatePositiveIntegerOption,
"OpenRestyProxyReadTimeout": validatePositiveIntegerOption,
"OpenRestyGzipMinLength": validatePositiveIntegerOption,
"OpenRestyGzipCompLevel": validateOpenRestyGzipCompLevel,
"OpenRestyEventsUse": validateOpenRestyEventsUse,
"OpenRestyResolvers": validateOpenRestyResolvers,
"OpenRestyEventsMultiAcceptEnabled": validateBooleanOption,
"OpenRestyWebsocketEnabled": validateBooleanOption,
"OpenRestyHTTP3Enabled": validateBooleanOption,
"OpenRestyProxyRequestBufferingEnabled": validateBooleanOption,
"OpenRestyProxyBufferingEnabled": validateBooleanOption,
"OpenRestyGzipEnabled": validateBooleanOption,
"OpenRestyCacheEnabled": validateBooleanOption,
"OpenRestyCacheLockEnabled": validateBooleanOption,
"OpenRestyProxyBuffers": validateOpenRestyProxyBuffers,
"OpenRestyLargeClientHeaderBuffers": validateOpenRestyProxyBuffers,
"OpenRestyProxyBufferSize": validateOpenRestySizeValue,
"OpenRestyProxyBusyBuffersSize": validateOpenRestySizeValue,
"OpenRestyCacheMaxSize": validateOpenRestySizeValue,
"OpenRestyClientMaxBodySize": validateOpenRestySizeValue,
"OpenRestyCachePath": validateOpenRestyCachePath,
"OpenRestyCacheLevels": validateOpenRestyCacheLevels,
"OpenRestyCacheInactive": validateOpenRestyDurationToken,
"OpenRestyCacheLockTimeout": validateOpenRestyDurationToken,
"OpenRestyCacheKeyTemplate": validateOpenRestyCacheKeyTemplate,
"OpenRestyCacheUseStale": validateOpenRestyCacheUseStale,
"OpenRestyMainConfigTemplate": validateOpenRestyMainConfigTemplate,
model.ConfigKeyOpenRestyDefaultServerReturnStatus: validateOpenRestyDefaultServerReturnStatus,
model.ConfigKeyOpenRestyWorkerProcesses: validateOpenRestyWorkerProcesses,
model.ConfigKeyOpenRestyWorkerConnections: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyWorkerRlimitNofile: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyKeepaliveTimeout: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyKeepaliveRequests: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyClientHeaderTimeout: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyClientBodyTimeout: validatePositiveIntegerOption,
model.ConfigKeyOpenRestySendTimeout: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyProxyConnectTimeout: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyProxySendTimeout: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyProxyReadTimeout: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyGzipMinLength: validatePositiveIntegerOption,
model.ConfigKeyOpenRestyGzipCompLevel: validateOpenRestyGzipCompLevel,
model.ConfigKeyOpenRestyEventsUse: validateOpenRestyEventsUse,
model.ConfigKeyOpenRestyResolvers: validateOpenRestyResolvers,
model.ConfigKeyOpenRestyEventsMultiAcceptEnabled: validateBooleanOption,
model.ConfigKeyOpenRestyWebsocketEnabled: validateBooleanOption,
model.ConfigKeyOpenRestyHTTP3Enabled: validateBooleanOption,
model.ConfigKeyOpenRestyProxyRequestBufferingEnabled: validateBooleanOption,
model.ConfigKeyOpenRestyProxyBufferingEnabled: validateBooleanOption,
model.ConfigKeyOpenRestyGzipEnabled: validateBooleanOption,
model.ConfigKeyOpenRestyCacheEnabled: validateBooleanOption,
model.ConfigKeyOpenRestyCacheLockEnabled: validateBooleanOption,
model.ConfigKeyOpenRestyProxyBuffers: validateOpenRestyProxyBuffers,
model.ConfigKeyOpenRestyLargeClientHeaderBuffers: validateOpenRestyProxyBuffers,
model.ConfigKeyOpenRestyProxyBufferSize: validateOpenRestySizeValue,
model.ConfigKeyOpenRestyProxyBusyBuffersSize: validateOpenRestySizeValue,
model.ConfigKeyOpenRestyCacheMaxSize: validateOpenRestySizeValue,
model.ConfigKeyOpenRestyClientMaxBodySize: validateOpenRestySizeValue,
model.ConfigKeyOpenRestyCachePath: validateOpenRestyCachePath,
model.ConfigKeyOpenRestyCacheLevels: validateOpenRestyCacheLevels,
model.ConfigKeyOpenRestyCacheInactive: validateOpenRestyDurationToken,
model.ConfigKeyOpenRestyCacheLockTimeout: validateOpenRestyDurationToken,
model.ConfigKeyOpenRestyCacheKeyTemplate: validateOpenRestyCacheKeyTemplate,
model.ConfigKeyOpenRestyCacheUseStale: validateOpenRestyCacheUseStale,
model.ConfigKeyOpenRestyMainConfigTemplate: validateOpenRestyMainConfigTemplate,
}
func validateOpenRestyOption(key, value string) error {
+41 -31
View File
@@ -4,6 +4,7 @@
package option
import (
"context"
"errors"
"fmt"
"regexp"
@@ -12,6 +13,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/apps/openflare/geoip"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const maxOpenRestyGzipCompLevel = 9
@@ -25,21 +27,25 @@ var (
const optionValueTrue = "true"
func buildOptionValidationState(options []model.OpenFlareOption) map[string]string {
model.OptionMapRWMutex.RLock()
state := make(map[string]string, len(model.OptionMap)+len(options))
for key, value := range model.OptionMap {
state[key] = value
}
model.OptionMapRWMutex.RUnlock()
func buildOptionValidationState(ctx context.Context, options []model.OpenFlareOption) map[string]string {
// 从 SystemConfig 读取所有业务配置构建状态
configs, err := repository.ListAdminSystemConfigs(ctx, "business")
state := make(map[string]string, len(configs)+len(options))
if err == nil {
for _, config := range configs {
state[config.Key] = config.Value
}
}
// 应用待验证的新值
for _, option := range options {
state[option.Key] = option.Value
}
return state
}
func validateOptionWithState(option model.OpenFlareOption, state map[string]string) error {
func validateOptionWithState(ctx context.Context, option model.OpenFlareOption, state map[string]string) error {
if err := validateOpenRestyOption(option.Key, option.Value); err != nil {
return err
@@ -53,7 +59,7 @@ func validateOptionWithState(option model.OpenFlareOption, state map[string]stri
if err := validateAgentOption(option.Key, option.Value); err != nil {
return err
}
return validateUptimeKumaOption(option.Key, option.Value, state)
return validateUptimeKumaOption(ctx, option.Key, option.Value, state)
}
func validatePositiveIntegerOption(key, value string) error {
@@ -74,7 +80,7 @@ func validateBooleanOption(key, value string) error {
}
func validateGeoIPOption(key, value string) error {
if key != "GeoIPProvider" {
if key != model.ConfigKeyGeoIPProvider {
return nil
}
if geoip.IsValidProvider(value) {
@@ -85,9 +91,9 @@ func validateGeoIPOption(key, value string) error {
func validateDatabaseCleanupOption(key, value string) error {
switch key {
case "DatabaseAutoCleanupEnabled":
case model.ConfigKeyDatabaseAutoCleanupEnabled:
return validateBooleanOption(key, value)
case "DatabaseAutoCleanupRetentionDays":
case model.ConfigKeyDatabaseAutoCleanupRetentionDays:
intValue, err := strconv.Atoi(value)
if err != nil || intValue < 1 {
return fmt.Errorf("%s 必须为大于等于 1 的整数天", key)
@@ -97,55 +103,59 @@ func validateDatabaseCleanupOption(key, value string) error {
}
func validateAgentOption(key, value string) error {
if key == "AgentWebsocketUpgradeEnabled" {
if key == model.ConfigKeyAgentWebsocketUpgradeEnabled {
return validateBooleanOption(key, strings.TrimSpace(value))
}
return nil
}
func validateUptimeKumaOption(key, value string, state map[string]string) error {
func validateUptimeKumaOption(ctx context.Context, key, value string, state map[string]string) error {
trimmed := strings.TrimSpace(value)
switch key {
case "UptimeKumaEnabled":
return validateUptimeKumaEnabled(key, trimmed, state)
case "UptimeKumaUsername":
case model.ConfigKeyUptimeKumaEnabled:
return validateUptimeKumaEnabled(ctx, key, trimmed, state)
case model.ConfigKeyUptimeKumaUsername:
return validateUptimeKumaUsername(trimmed, state)
case "UptimeKumaUrl":
case model.ConfigKeyUptimeKumaURL:
return validateUptimeKumaURL(trimmed)
case "UptimeKumaMonitorScope":
case model.ConfigKeyUptimeKumaMonitorScope:
return validateUptimeKumaMonitorScope(trimmed)
case "UptimeKumaSyncInterval", "UptimeKumaInterval", "UptimeKumaRetryInterval", "UptimeKumaTimeout":
case model.ConfigKeyUptimeKumaSyncInterval, model.ConfigKeyUptimeKumaInterval, model.ConfigKeyUptimeKumaRetryInterval, model.ConfigKeyUptimeKumaTimeout:
return validatePositiveIntegerOption(key, trimmed)
case "UptimeKumaRetry":
case model.ConfigKeyUptimeKumaRetry:
return validateUptimeKumaRetry(key, trimmed)
}
return nil
}
func validateUptimeKumaEnabled(key, trimmed string, state map[string]string) error {
func validateUptimeKumaEnabled(ctx context.Context, key, trimmed string, state map[string]string) error {
if err := validateBooleanOption(key, trimmed); err != nil {
return err
}
if trimmed != optionValueTrue {
return nil
}
url := strings.TrimSpace(state["UptimeKumaUrl"])
username := strings.TrimSpace(state["UptimeKumaUsername"])
password := strings.TrimSpace(state["UptimeKumaPassword"])
url := strings.TrimSpace(state[model.ConfigKeyUptimeKumaURL])
username := strings.TrimSpace(state[model.ConfigKeyUptimeKumaUsername])
password := strings.TrimSpace(state[model.ConfigKeyUptimeKumaPassword])
if url == "" {
return fmt.Errorf("启用 Uptime Kuma 时地址不能为空")
}
if username == "" {
return fmt.Errorf("启用 Uptime Kuma 时用户名不能为空")
}
if password == "" && model.UptimeKumaPassword == "" {
return fmt.Errorf("启用 Uptime Kuma 时密码不能为空")
// 如果待验证的密码为空,且当前配置中也没有密码,则报错
if password == "" {
existingPwd, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUptimeKumaPassword)
if strings.TrimSpace(existingPwd.Value) == "" {
return fmt.Errorf("启用 Uptime Kuma 时密码不能为空")
}
}
return nil
}
func validateUptimeKumaUsername(trimmed string, state map[string]string) error {
if trimmed == "" && state["UptimeKumaEnabled"] == optionValueTrue {
if trimmed == "" && state[model.ConfigKeyUptimeKumaEnabled] == optionValueTrue {
return fmt.Errorf("启用 Uptime Kuma 时用户名不能为空")
}
return nil
@@ -173,17 +183,17 @@ func validateUptimeKumaRetry(key, trimmed string) error {
return nil
}
func validateOptions(options []model.OpenFlareOption) error {
func validateOptions(ctx context.Context, options []model.OpenFlareOption) error {
if len(options) == 0 {
return errors.New(errInvalidParams)
}
state := buildOptionValidationState(options)
state := buildOptionValidationState(ctx, options)
for _, option := range options {
if strings.TrimSpace(option.Key) == "" {
return errors.New(errInvalidParams)
}
if err := validateOptionWithState(option, state); err != nil {
if err := validateOptionWithState(ctx, option, state); err != nil {
return err
}
}
+19 -7
View File
@@ -15,6 +15,9 @@ import (
const (
relayStatusUnhealthy = "unhealthy"
releaseChannelStable = "stable"
defaultAgentHeartbeatInterval = 10000 // 默认心跳间隔 10 秒(毫秒)
defaultAgentUpdateRepo = "Rain-kl/OpenFlare"
)
func normalizeRelayStatus(status string) string {
@@ -104,10 +107,7 @@ func buildRelayConfig(ctx context.Context, node *model.OpenFlareNode) *Config {
}
// BuildSettings returns runtime settings shared by relay and flared clients.
func BuildSettings(node *model.OpenFlareNode, updateNow bool, updateChannel, updateTag string) *Settings {
model.OptionMapRWMutex.RLock()
defer model.OptionMapRWMutex.RUnlock()
func BuildSettings(ctx context.Context, node *model.OpenFlareNode, updateNow bool, updateChannel, updateTag string) *Settings {
autoUpdate := false
if node != nil {
autoUpdate = node.AutoUpdateEnabled
@@ -115,11 +115,23 @@ func BuildSettings(node *model.OpenFlareNode, updateNow bool, updateChannel, upd
if strings.TrimSpace(updateChannel) == "" {
updateChannel = releaseChannelStable
}
// 从 SystemConfig 读取配置,使用默认值作为降级
heartbeatInterval, _ := repository.GetIntByKey(ctx, model.ConfigKeyAgentHeartbeatInterval)
if heartbeatInterval <= 0 {
heartbeatInterval = defaultAgentHeartbeatInterval
}
wsUpgradeEnabled, _ := repository.GetBoolByKey(ctx, model.ConfigKeyAgentWebsocketUpgradeEnabled)
updateRepo, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyAgentUpdateRepo)
if strings.TrimSpace(updateRepo.Value) == "" {
updateRepo.Value = defaultAgentUpdateRepo
}
return &Settings{
HeartbeatInterval: model.AgentHeartbeatInterval,
WebsocketUpgradeEnabled: model.AgentWebsocketUpgradeEnabled,
HeartbeatInterval: heartbeatInterval,
WebsocketUpgradeEnabled: wsUpgradeEnabled,
AutoUpdate: autoUpdate,
UpdateRepo: model.AgentUpdateRepo,
UpdateRepo: updateRepo.Value,
UpdateNow: updateNow,
UpdateChannel: updateChannel,
UpdateTag: strings.TrimSpace(updateTag),
+1 -1
View File
@@ -99,7 +99,7 @@ func Heartbeat(ctx context.Context, node *model.OpenFlareNode, payload Heartbeat
return &HeartbeatResponse{
RelayConfig: buildRelayConfig(ctx, node),
RelaySettings: BuildSettings(node, updateNow, updateChannel, updateTag),
RelaySettings: BuildSettings(ctx, node, updateNow, updateChannel, updateTag),
}, nil
}
+1 -4
View File
@@ -10,7 +10,6 @@ import (
"time"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/agent"
"github.com/Rain-kl/Wavelet/internal/apps/openflare/option"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/glebarez/sqlite"
@@ -28,7 +27,7 @@ func setupRelayTestDB(t *testing.T) func() {
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(
&model.OpenFlareNode{},
&model.OpenFlareOption{},
&model.SystemConfig{},
&model.OpenFlareNodeSystemProfile{},
&model.OpenFlareMetricSnapshot{},
&model.OpenFlareHealthEvent{},
@@ -36,14 +35,12 @@ func setupRelayTestDB(t *testing.T) func() {
))
db.SetDB(sqliteDB)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore())
return func() {
resetObservabilityStore()
db.SetDB(nil)
option.ResetInitializationForTest()
agent.ResetAuthCacheForTest()
}
}
@@ -11,6 +11,7 @@ import (
"time"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const (
@@ -91,14 +92,20 @@ func CleanupDatabaseObservability(ctx context.Context, input DatabaseCleanupInpu
// RunDatabaseAutoCleanupOnce runs retention-based cleanup for all observability targets.
func RunDatabaseAutoCleanupOnce(ctx context.Context, now time.Time) (*DatabaseAutoCleanupSummary, error) {
if !model.DatabaseAutoCleanupEnabled {
enabled, err := repository.GetBoolByKey(ctx, model.ConfigKeyDatabaseAutoCleanupEnabled)
if err != nil {
return nil, fmt.Errorf("failed to read database_auto_cleanup_enabled: %w", err)
}
if !enabled {
return nil, nil
}
if model.DatabaseAutoCleanupRetentionDays < 1 {
return nil, fmt.Errorf("database auto cleanup retention_days must be at least 1")
retentionDays, err := repository.GetIntByKey(ctx, model.ConfigKeyDatabaseAutoCleanupRetentionDays)
if err != nil || retentionDays <= 0 {
// Use default value 30 if config read fails or value is invalid
retentionDays = 30
}
retentionDays := model.DatabaseAutoCleanupRetentionDays
results := make([]DatabaseCleanupResult, 0, len(databaseCleanupTargets))
for _, target := range []string{
DatabaseCleanupTargetAccessLogs,
@@ -10,6 +10,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -23,6 +24,7 @@ func setupDatabaseCleanupTestDB(t *testing.T) context.Context {
DisableForeignKeyConstraintWhenMigrating: true,
})
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(&model.SystemConfig{}))
db.SetDB(sqliteDB)
resetAccessLogStore := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore())
resetObservabilityStore := model.SetObservabilityStoreForTest(model.NewMemoryObservabilityStore())
@@ -123,14 +125,8 @@ func TestRunDatabaseAutoCleanupOnceDeletesAllObservabilityTargets(t *testing.T)
RequestCount: 15,
}))
previousEnabled := model.DatabaseAutoCleanupEnabled
previousRetentionDays := model.DatabaseAutoCleanupRetentionDays
model.DatabaseAutoCleanupEnabled = true
model.DatabaseAutoCleanupRetentionDays = 1
t.Cleanup(func() {
model.DatabaseAutoCleanupEnabled = previousEnabled
model.DatabaseAutoCleanupRetentionDays = previousRetentionDays
})
require.NoError(t, repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyDatabaseAutoCleanupEnabled, "true"))
require.NoError(t, repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyDatabaseAutoCleanupRetentionDays, "1"))
summary, err := RunDatabaseAutoCleanupOnce(ctx, now)
require.NoError(t, err)
+86 -25
View File
@@ -12,15 +12,70 @@ import (
"github.com/Rain-kl/Wavelet/internal/apps/openflare/routeidentity"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const uptimeKumaTagOpenFlare = "OpenFlare"
var isSyncing atomic.Bool
// kumaConfig 封装 UptimeKuma 配置
type kumaConfig struct {
URL string
Username string
Password string
MonitorScope string
SelectedSites string
Interval int
Retry int
RetryInterval int
Timeout int
}
// loadKumaConfig 从 SystemConfig 加载 UptimeKuma 配置
func loadKumaConfig(ctx context.Context) (*kumaConfig, error) {
url, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUptimeKumaURL)
username, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUptimeKumaUsername)
password, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUptimeKumaPassword)
scope, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUptimeKumaMonitorScope)
selected, _ := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUptimeKumaSelectedSites)
interval, _ := repository.GetIntByKey(ctx, model.ConfigKeyUptimeKumaInterval)
if interval <= 0 {
interval = 60
}
retry, _ := repository.GetIntByKey(ctx, model.ConfigKeyUptimeKumaRetry)
retryInterval, _ := repository.GetIntByKey(ctx, model.ConfigKeyUptimeKumaRetryInterval)
if retryInterval <= 0 {
retryInterval = 60
}
timeout, _ := repository.GetIntByKey(ctx, model.ConfigKeyUptimeKumaTimeout)
if timeout <= 0 {
timeout = 48
}
if scope.Value == "" {
scope.Value = "all"
}
return &kumaConfig{
URL: strings.TrimSpace(url.Value),
Username: strings.TrimSpace(username.Value),
Password: strings.TrimSpace(password.Value),
MonitorScope: scope.Value,
SelectedSites: selected.Value,
Interval: interval,
Retry: retry,
RetryInterval: retryInterval,
Timeout: timeout,
}, nil
}
// SyncToUptimeKuma synchronizes enabled proxy routes to Uptime Kuma monitors.
func SyncToUptimeKuma(ctx context.Context) error {
if !model.UptimeKumaEnabled {
// 检查是否启用
enabled, _ := repository.GetBoolByKey(ctx, model.ConfigKeyUptimeKumaEnabled)
if !enabled {
return fmt.Errorf("uptime Kuma integration is disabled")
}
@@ -29,15 +84,21 @@ func SyncToUptimeKuma(ctx context.Context) error {
}
defer isSyncing.Store(false)
kumaURL, kumaUsername, kumaPassword, err := validateUptimeKumaConfig()
// 加载配置
config, err := loadKumaConfig(ctx)
if err != nil {
return err
}
// 验证配置
if err := validateKumaConfig(config); err != nil {
return err
}
slog.Info("Starting Uptime Kuma sync process",
"url", kumaURL,
"username", kumaUsername,
"scope", model.UptimeKumaMonitorScope,
"url", config.URL,
"username", config.Username,
"scope", config.MonitorScope,
)
allRoutes, err := model.ListProxyRoutes(ctx)
@@ -45,12 +106,12 @@ func SyncToUptimeKuma(ctx context.Context) error {
return fmt.Errorf("failed to list local proxy routes: %w", err)
}
expectedRoutes, err := filterExpectedRoutes(allRoutes)
expectedRoutes, err := filterExpectedRoutes(allRoutes, config)
if err != nil {
return err
}
client, err := connectAndLoginUptimeKuma(kumaURL, kumaUsername, kumaPassword)
client, err := connectAndLoginUptimeKuma(config.URL, config.Username, config.Password)
if err != nil {
return err
}
@@ -62,16 +123,16 @@ func SyncToUptimeKuma(ctx context.Context) error {
}
existingOpenFlareMonitors := filterOpenFlareMonitors(client.GetMonitorList(), openFlareTagID)
expectedSitesMap := syncRouteMonitors(client, expectedRoutes, existingOpenFlareMonitors, openFlareTagID)
expectedSitesMap := syncRouteMonitors(client, expectedRoutes, existingOpenFlareMonitors, openFlareTagID, config)
removeStaleMonitors(client, existingOpenFlareMonitors, expectedSitesMap)
return nil
}
func filterExpectedRoutes(allRoutes []*model.ProxyRoute) ([]*model.ProxyRoute, error) {
scope := model.UptimeKumaMonitorScope
func filterExpectedRoutes(allRoutes []*model.ProxyRoute, config *kumaConfig) ([]*model.ProxyRoute, error) {
scope := config.MonitorScope
if scope == "selected" {
selectedList := strings.Split(model.UptimeKumaSelectedSites, ",")
selectedList := strings.Split(config.SelectedSites, ",")
selectedMap := make(map[string]bool)
for _, name := range selectedList {
trimmedName := strings.TrimSpace(name)
@@ -175,15 +236,15 @@ func routeMonitorURL(route *model.ProxyRoute) (string, error) {
return "http://" + domain, nil
}
func monitorPayload(id int, name, targetURL string) map[string]any {
func monitorPayload(id int, name, targetURL string, config *kumaConfig) map[string]any {
payload := map[string]any{
"type": "http",
"name": name,
"url": targetURL,
"interval": model.UptimeKumaInterval,
"maxretries": model.UptimeKumaRetry,
"retryInterval": model.UptimeKumaRetryInterval,
"timeout": model.UptimeKumaTimeout,
"interval": config.Interval,
"maxretries": config.Retry,
"retryInterval": config.RetryInterval,
"timeout": config.Timeout,
"active": true,
"resendInterval": 0,
"expiryNotification": false,
@@ -198,17 +259,17 @@ func monitorPayload(id int, name, targetURL string) map[string]any {
return payload
}
func monitorNeedsUpdate(existing Monitor, targetURL string) bool {
func monitorNeedsUpdate(existing Monitor, targetURL string, config *kumaConfig) bool {
return existing.URL != targetURL ||
existing.Interval != model.UptimeKumaInterval ||
existing.MaxRetries != model.UptimeKumaRetry ||
existing.RetryInterval != model.UptimeKumaRetryInterval ||
existing.Timeout != model.UptimeKumaTimeout
existing.Interval != config.Interval ||
existing.MaxRetries != config.Retry ||
existing.RetryInterval != config.RetryInterval ||
existing.Timeout != config.Timeout
}
func createMonitor(client *SocketIOClient, siteName, targetURL string, openFlareTagID int) error {
func createMonitor(client *SocketIOClient, siteName, targetURL string, openFlareTagID int, config *kumaConfig) error {
slog.Info("Creating monitor in Uptime Kuma", "name", siteName, "url", targetURL)
addAck, err := client.Emit("add", monitorPayload(0, siteName, targetURL))
addAck, err := client.Emit("add", monitorPayload(0, siteName, targetURL, config))
if err != nil {
return err
}
@@ -238,9 +299,9 @@ func createMonitor(client *SocketIOClient, siteName, targetURL string, openFlare
return nil
}
func updateMonitor(client *SocketIOClient, monitorID int, siteName, targetURL string) error {
func updateMonitor(client *SocketIOClient, monitorID int, siteName, targetURL string, config *kumaConfig) error {
slog.Info("Updating monitor in Uptime Kuma due to settings mismatch", "name", siteName)
editAck, err := client.Emit("editMonitor", monitorPayload(monitorID, siteName, targetURL))
editAck, err := client.Emit("editMonitor", monitorPayload(monitorID, siteName, targetURL, config))
if err != nil {
return err
}
@@ -14,17 +14,18 @@ import (
const monitorListWaitTimeout = 5 * time.Second
func validateUptimeKumaConfig() (string, string, string, error) {
kumaURL := strings.TrimSpace(model.UptimeKumaURL)
kumaUsername := strings.TrimSpace(model.UptimeKumaUsername)
kumaPassword := strings.TrimSpace(model.UptimeKumaPassword)
if kumaURL == "" || kumaUsername == "" || kumaPassword == "" {
return kumaURL, kumaUsername, kumaPassword, fmt.Errorf(
"uptime Kuma URL, username, or password is not configured (URL: %q, Username: %q, PasswordLength: %d)",
kumaURL, kumaUsername, len(kumaPassword),
)
// validateKumaConfig 验证 kumaConfig 配置完整性
func validateKumaConfig(config *kumaConfig) error {
if strings.TrimSpace(config.URL) == "" {
return fmt.Errorf("uptime Kuma URL is not configured")
}
return kumaURL, kumaUsername, kumaPassword, nil
if strings.TrimSpace(config.Username) == "" {
return fmt.Errorf("uptime Kuma username is not configured")
}
if strings.TrimSpace(config.Password) == "" {
return fmt.Errorf("uptime Kuma password is not configured")
}
return nil
}
func connectAndLoginUptimeKuma(kumaURL, kumaUsername, kumaPassword string) (*SocketIOClient, error) {
@@ -68,7 +69,7 @@ func connectAndLoginUptimeKuma(kumaURL, kumaUsername, kumaPassword string) (*Soc
return client, nil
}
func syncRouteMonitors(client *SocketIOClient, expectedRoutes []*model.ProxyRoute, existingMonitors map[string]Monitor, openFlareTagID int) map[string]bool {
func syncRouteMonitors(client *SocketIOClient, expectedRoutes []*model.ProxyRoute, existingMonitors map[string]Monitor, openFlareTagID int, config *kumaConfig) map[string]bool {
expectedSitesMap := make(map[string]bool, len(expectedRoutes))
for _, route := range expectedRoutes {
expectedSitesMap[route.SiteName] = true
@@ -80,13 +81,13 @@ func syncRouteMonitors(client *SocketIOClient, expectedRoutes []*model.ProxyRout
existing, exists := existingMonitors[route.SiteName]
if !exists {
if err := createMonitor(client, route.SiteName, targetURL, openFlareTagID); err != nil {
if err := createMonitor(client, route.SiteName, targetURL, openFlareTagID, config); err != nil {
slog.Error("Failed to add monitor to Uptime Kuma", "name", route.SiteName, "error", err)
}
continue
}
if monitorNeedsUpdate(existing, targetURL) {
if err := updateMonitor(client, existing.ID, route.SiteName, targetURL); err != nil {
if monitorNeedsUpdate(existing, targetURL, config) {
if err := updateMonitor(client, existing.ID, route.SiteName, targetURL, config); err != nil {
slog.Error("Failed to edit monitor in Uptime Kuma", "name", route.SiteName, "error", err)
}
}
+56 -44
View File
@@ -17,6 +17,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -121,7 +122,7 @@ func setupSyncTestDB(t *testing.T) func() {
DisableForeignKeyConstraintWhenMigrating: true,
})
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(&model.ProxyRoute{}))
require.NoError(t, sqliteDB.AutoMigrate(&model.ProxyRoute{}, &model.SystemConfig{}))
db.SetDB(sqliteDB)
return func() {
@@ -129,41 +130,50 @@ func setupSyncTestDB(t *testing.T) func() {
}
}
func backupUptimeKumaConfig() func() {
oldEnabled := model.UptimeKumaEnabled
oldURL := model.UptimeKumaURL
oldUsername := model.UptimeKumaUsername
oldPassword := model.UptimeKumaPassword
oldScope := model.UptimeKumaMonitorScope
oldSelected := model.UptimeKumaSelectedSites
oldInterval := model.UptimeKumaInterval
oldRetry := model.UptimeKumaRetry
oldRetryInterval := model.UptimeKumaRetryInterval
oldTimeout := model.UptimeKumaTimeout
func backupUptimeKumaConfig(ctx context.Context) func() {
// 备份所有 UptimeKuma 相关配置
configs := []string{
model.ConfigKeyUptimeKumaEnabled,
model.ConfigKeyUptimeKumaURL,
model.ConfigKeyUptimeKumaUsername,
model.ConfigKeyUptimeKumaPassword,
model.ConfigKeyUptimeKumaMonitorScope,
model.ConfigKeyUptimeKumaSelectedSites,
model.ConfigKeyUptimeKumaInterval,
model.ConfigKeyUptimeKumaRetry,
model.ConfigKeyUptimeKumaRetryInterval,
model.ConfigKeyUptimeKumaTimeout,
}
oldValues := make(map[string]string)
for _, key := range configs {
config, _ := repository.GetSystemConfigByKey(ctx, key)
oldValues[key] = config.Value
}
return func() {
model.UptimeKumaEnabled = oldEnabled
model.UptimeKumaURL = oldURL
model.UptimeKumaUsername = oldUsername
model.UptimeKumaPassword = oldPassword
model.UptimeKumaMonitorScope = oldScope
model.UptimeKumaSelectedSites = oldSelected
model.UptimeKumaInterval = oldInterval
model.UptimeKumaRetry = oldRetry
model.UptimeKumaRetryInterval = oldRetryInterval
model.UptimeKumaTimeout = oldTimeout
// 恢复所有配置
for key, value := range oldValues {
_ = db.DB(ctx).Model(&model.SystemConfig{}).Where("key = ?", key).Update("value", value).Error
}
}
}
// setTestConfig 设置测试配置的辅助函数(不存在则创建)
func setTestConfig(ctx context.Context, key, value string) {
_ = repository.SaveOrUpdateSystemConfig(ctx, key, value)
}
func TestSyncToUptimeKumaDisabled(t *testing.T) {
cleanup := setupSyncTestDB(t)
defer cleanup()
restore := backupUptimeKumaConfig()
ctx := context.Background()
restore := backupUptimeKumaConfig(ctx)
defer restore()
model.UptimeKumaEnabled = false
setTestConfig(ctx, model.ConfigKeyUptimeKumaEnabled, "false")
err := SyncToUptimeKuma(context.Background())
err := SyncToUptimeKuma(ctx)
require.Error(t, err)
assert.Contains(t, err.Error(), "disabled")
}
@@ -171,9 +181,9 @@ func TestSyncToUptimeKumaDisabled(t *testing.T) {
func TestSyncToUptimeKumaSuccess(t *testing.T) {
cleanup := setupSyncTestDB(t)
defer cleanup()
restore := backupUptimeKumaConfig()
defer restore()
ctx := context.Background()
restore := backupUptimeKumaConfig(ctx)
defer restore()
require.NoError(t, db.DB(ctx).Where("1 = 1").Delete(&model.ProxyRoute{}).Error)
@@ -227,15 +237,16 @@ func TestSyncToUptimeKumaSuccess(t *testing.T) {
server := httptest.NewServer(mockSrv)
defer server.Close()
model.UptimeKumaEnabled = true
model.UptimeKumaURL = server.URL
model.UptimeKumaUsername = "admin"
model.UptimeKumaPassword = "password"
model.UptimeKumaMonitorScope = "all"
model.UptimeKumaInterval = 60
model.UptimeKumaRetry = 0
model.UptimeKumaRetryInterval = 60
model.UptimeKumaTimeout = 48
// 设置测试配置
setTestConfig(ctx, model.ConfigKeyUptimeKumaEnabled, "true")
setTestConfig(ctx, model.ConfigKeyUptimeKumaURL, server.URL)
setTestConfig(ctx, model.ConfigKeyUptimeKumaUsername, "admin")
setTestConfig(ctx, model.ConfigKeyUptimeKumaPassword, "password")
setTestConfig(ctx, model.ConfigKeyUptimeKumaMonitorScope, "all")
setTestConfig(ctx, model.ConfigKeyUptimeKumaInterval, "60")
setTestConfig(ctx, model.ConfigKeyUptimeKumaRetry, "0")
setTestConfig(ctx, model.ConfigKeyUptimeKumaRetryInterval, "60")
setTestConfig(ctx, model.ConfigKeyUptimeKumaTimeout, "48")
require.NoError(t, SyncToUptimeKuma(ctx))
@@ -282,9 +293,9 @@ func TestSyncToUptimeKumaSuccess(t *testing.T) {
func TestSyncToUptimeKumaSelectedScope(t *testing.T) {
cleanup := setupSyncTestDB(t)
defer cleanup()
restore := backupUptimeKumaConfig()
defer restore()
ctx := context.Background()
restore := backupUptimeKumaConfig(ctx)
defer restore()
require.NoError(t, db.DB(ctx).Where("1 = 1").Delete(&model.ProxyRoute{}).Error)
@@ -312,12 +323,13 @@ func TestSyncToUptimeKumaSelectedScope(t *testing.T) {
server := httptest.NewServer(mockSrv)
defer server.Close()
model.UptimeKumaEnabled = true
model.UptimeKumaURL = server.URL
model.UptimeKumaUsername = "admin"
model.UptimeKumaPassword = "password"
model.UptimeKumaMonitorScope = "selected"
model.UptimeKumaSelectedSites = "site-a"
// 设置测试配置
setTestConfig(ctx, model.ConfigKeyUptimeKumaEnabled, "true")
setTestConfig(ctx, model.ConfigKeyUptimeKumaURL, server.URL)
setTestConfig(ctx, model.ConfigKeyUptimeKumaUsername, "admin")
setTestConfig(ctx, model.ConfigKeyUptimeKumaPassword, "password")
setTestConfig(ctx, model.ConfigKeyUptimeKumaMonitorScope, "selected")
setTestConfig(ctx, model.ConfigKeyUptimeKumaSelectedSites, "site-a")
require.NoError(t, SyncToUptimeKuma(ctx))