From f671a96d8c8217fe5a6f31eb8f2e9c9519185ec9 Mon Sep 17 00:00:00 2001 From: ryan Date: Wed, 3 Jun 2026 11:14:19 +0800 Subject: [PATCH] =?UTF-8?q?[=E6=96=B0=E5=A2=9E]=20=E5=AF=B9=E6=8E=A5=20Upt?= =?UTF-8?q?ime=20Kuma?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/guildline/development-constraints.md | 4 +- docs/reference/configuration.md | 11 + openflare_server/common/constants.go | 13 + openflare_server/controller/option.go | 90 ++-- openflare_server/controller/option_test.go | 56 ++- openflare_server/controller/uptimekuma.go | 24 + openflare_server/job/cron.go | 8 + openflare_server/job/uptimekuma.go | 43 ++ openflare_server/model/option.go | 49 ++ openflare_server/router/api-router.go | 5 + .../router/api_uptimekuma_test.go | 454 ++++++++++++++++++ openflare_server/service/uptimekuma.go | 302 ++++++++++++ openflare_server/utils/uptimekuma/client.go | 373 ++++++++++++++ .../web/features/settings/api/settings.ts | 6 + .../settings/components/settings-page.tsx | 291 +++++++++++ .../settings/components/uptimekuma-modal.tsx | 211 ++++++++ 16 files changed, 1906 insertions(+), 34 deletions(-) create mode 100644 openflare_server/controller/uptimekuma.go create mode 100644 openflare_server/job/uptimekuma.go create mode 100644 openflare_server/router/api_uptimekuma_test.go create mode 100644 openflare_server/service/uptimekuma.go create mode 100644 openflare_server/utils/uptimekuma/client.go create mode 100644 openflare_server/web/features/settings/components/uptimekuma-modal.tsx diff --git a/docs/guildline/development-constraints.md b/docs/guildline/development-constraints.md index 609279f7..48ef3c9b 100644 --- a/docs/guildline/development-constraints.md +++ b/docs/guildline/development-constraints.md @@ -49,7 +49,9 @@ Frontend: 各组件和模块(Server、Agent、Frontend)的物理目录分层职责详见 [仓库结构](../design/repository.md)。在此结构下,开发必须遵守以下核心分层规则: -* **Server 开发规则**:禁止在 `controller/` 堆积业务逻辑,禁止在 `middleware/` 实现业务流程,禁止为简单需求新增平台层抽象。 +* **Server 开发规则**: + * 禁止在 `controller/` 堆积业务逻辑,禁止在 `middleware/` 实现业务流程,禁止为简单需求新增平台层抽象。 + * **定时任务开发规则**:禁止将不同业务模块(如 Uptime Kuma 整合、WAF IP 同步等)的定时任务具体执行逻辑与状态堆积在单个 `cron.go` 文件中。各模块对应的定时任务结构体和运行逻辑必须在独立的 Go 文件中定义,`cron.go` 只允许承担统一注册、初始化与调度器启停的职责。 * **Agent 开发规则**:每个模块职责单一,外部命令调用集中封装,状态落盘与配置落盘分离。 * **Frontend 开发规则**:页面文件只负责获取路由参数、组织页面结构、调用 feature 组件;不应手写复杂 API 细节、复杂表单校验逻辑或维护大量彼此耦合的局部状态。 diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index 689632a9..f60f5a1c 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -95,6 +95,17 @@ go run . --port 3000 --log-dir ./logs | `GlobalApiRateLimitNum` / `GlobalApiRateLimitDuration` | 全局 API 限流次数 / 时间窗口 | `300` / `180` | | `GlobalWebRateLimitNum` / `GlobalWebRateLimitDuration` | 全局 Web 限流次数 / 时间窗口 | `300` / `180` | | `CriticalRateLimitNum` / `CriticalRateLimitDuration` | 敏感接口限流次数 / 时间窗口 | `100` / `1200` | +| `UptimeKumaEnabled` | 是否启用 Uptime Kuma 自动同步 | `false` | +| `UptimeKumaUrl` | Uptime Kuma 实例地址 | 空 | +| `UptimeKumaUsername` | Uptime Kuma 登录用户名 | 空 | +| `UptimeKumaPassword` | Uptime Kuma 登录密码(写专,接口不回显) | 空 | +| `UptimeKumaMonitorScope` | 监控范围,支持 `all` (全部站点) 或 `selected` (选择站点) | `all` | +| `UptimeKumaSelectedSites` | 已选择监控站点的名称列表(英文逗号分隔) | 空 | +| `UptimeKumaSyncInterval` | 自动差分同步间隔(分钟) | `5` | +| `UptimeKumaInterval` | 监控心跳检测频率(秒) | `60` | +| `UptimeKumaRetry` | 监控最大重试次数 | `0` | +| `UptimeKumaRetryInterval` | 监控重试间隔时间(秒) | `60` | +| `UptimeKumaTimeout` | 监控请求超时断开时间(秒) | `48` | 说明: diff --git a/openflare_server/common/constants.go b/openflare_server/common/constants.go index 51b55660..25c22e9b 100644 --- a/openflare_server/common/constants.go +++ b/openflare_server/common/constants.go @@ -57,6 +57,19 @@ var GeoIPProvider = "ipinfo" var DatabaseAutoCleanupEnabled = false var DatabaseAutoCleanupRetentionDays = 30 +// Uptime Kuma integration settings +var UptimeKumaEnabled = false +var UptimeKumaUrl = "" +var UptimeKumaUsername = "" +var UptimeKumaPassword = "" +var UptimeKumaMonitorScope = "all" // "all" or "selected" +var UptimeKumaSelectedSites = "" // Comma-separated list of site names +var UptimeKumaSyncInterval = 5 // minutes +var UptimeKumaInterval = 60 // seconds +var UptimeKumaRetry = 0 +var UptimeKumaRetryInterval = 60 // seconds +var UptimeKumaTimeout = 48 // seconds + // V5 OpenResty performance settings (hot-reloadable via Option table) var OpenRestyDefaultServerReturnStatus = 421 diff --git a/openflare_server/controller/option.go b/openflare_server/controller/option.go index b178a0d4..648c10d1 100644 --- a/openflare_server/controller/option.go +++ b/openflare_server/controller/option.go @@ -100,6 +100,56 @@ func validateAgentOption(key string, value string) error { } } +func validateUptimeKumaOption(key string, value string, state map[string]string) error { + trimmed := strings.TrimSpace(value) + switch key { + case "UptimeKumaEnabled": + if err := validateBooleanOption(key, trimmed); err != nil { + return err + } + if trimmed == "true" { + url := strings.TrimSpace(state["UptimeKumaUrl"]) + username := strings.TrimSpace(state["UptimeKumaUsername"]) + password := strings.TrimSpace(state["UptimeKumaPassword"]) + if url == "" { + return fmt.Errorf("启用 Uptime Kuma 时地址不能为空") + } + if username == "" { + return fmt.Errorf("启用 Uptime Kuma 时用户名不能为空") + } + if password == "" && common.UptimeKumaPassword == "" { + return fmt.Errorf("启用 Uptime Kuma 时密码不能为空") + } + } + case "UptimeKumaUsername": + if trimmed == "" && state["UptimeKumaEnabled"] == "true" { + return fmt.Errorf("启用 Uptime Kuma 时用户名不能为空") + } + case "UptimeKumaPassword": + // No specific format checks needed + case "UptimeKumaUrl": + if trimmed != "" { + if !strings.HasPrefix(trimmed, "http://") && !strings.HasPrefix(trimmed, "https://") { + return fmt.Errorf("Uptime Kuma 地址必须以 http:// 或 https:// 开头") + } + } + case "UptimeKumaMonitorScope": + if trimmed != "all" && trimmed != "selected" { + return fmt.Errorf("监控范围必须为全部站点 (all) 或选择站点 (selected)") + } + case "UptimeKumaSyncInterval", "UptimeKumaInterval", "UptimeKumaRetryInterval", "UptimeKumaTimeout": + if err := validatePositiveIntegerOption(key, trimmed); err != nil { + return err + } + case "UptimeKumaRetry": + intValue, err := strconv.Atoi(trimmed) + if err != nil || intValue < 0 { + return fmt.Errorf("%s 必须为大于等于 0 的整数", key) + } + } + return nil +} + func validateOpenRestyOption(key string, value string) error { trimmed := strings.TrimSpace(value) @@ -263,6 +313,9 @@ func validateOptionWithState(option model.Option, state map[string]string) error if err := validateAgentOption(option.Key, option.Value); err != nil { return err } + if err := validateUptimeKumaOption(option.Key, option.Value, state); err != nil { + return err + } return nil } @@ -292,9 +345,9 @@ func updateOptions(options []model.Option) error { // @Router /api/option/ [get] func GetOptions(c *gin.Context) { var options []*model.Option - common.OptionMapRWMutex.Lock() + common.OptionMapRWMutex.RLock() for k, v := range common.OptionMap { - if strings.Contains(k, "Token") || strings.Contains(k, "Secret") { + if strings.Contains(k, "Token") || strings.Contains(k, "Secret") || strings.Contains(k, "Password") { continue } options = append(options, &model.Option{ @@ -302,7 +355,7 @@ func GetOptions(c *gin.Context) { Value: utils.Interface2String(v), }) } - common.OptionMapRWMutex.Unlock() + common.OptionMapRWMutex.RUnlock() respondSuccess(c, options) } @@ -320,35 +373,8 @@ func UpdateOption(c *gin.Context) { if !bindJSON(c, &option) { return } - switch option.Key { - case "GitHubOAuthEnabled": - if option.Value == "true" && common.GitHubClientId == "" { - respondFailure(c, "无法启用 GitHub OAuth,请先填入 GitHub Client ID 以及 GitHub Client Secret!") - return - } - case "WeChatAuthEnabled": - if option.Value == "true" && common.WeChatServerAddress == "" { - respondFailure(c, "无法启用微信登录,请先填入微信登录相关配置信息!") - return - } - } - if err := validateRateLimitOption(option.Key, option.Value); err != nil { - respondFailure(c, err.Error()) - return - } - if err := validateOpenRestyOption(option.Key, option.Value); err != nil { - respondFailure(c, err.Error()) - return - } - if err := validateGeoIPOption(option.Key, option.Value); err != nil { - respondFailure(c, err.Error()) - return - } - if err := validateDatabaseCleanupOption(option.Key, option.Value); err != nil { - respondFailure(c, err.Error()) - return - } - if err := validateAgentOption(option.Key, option.Value); err != nil { + state := buildOptionValidationState([]model.Option{option}) + if err := validateOptionWithState(option, state); err != nil { respondFailure(c, err.Error()) return } diff --git a/openflare_server/controller/option_test.go b/openflare_server/controller/option_test.go index a8b9d17b..92165f48 100644 --- a/openflare_server/controller/option_test.go +++ b/openflare_server/controller/option_test.go @@ -1,6 +1,8 @@ package controller -import "testing" +import ( + "testing" +) func TestValidateOpenRestyOption(t *testing.T) { testCases := []struct { @@ -63,3 +65,55 @@ func TestValidateAgentOption(t *testing.T) { t.Fatal("expected websocket upgrade option to reject non-boolean value") } } + +func TestValidateUptimeKumaOption(t *testing.T) { + state := map[string]string{ + "UptimeKumaUrl": "http://localhost:3001", + "UptimeKumaUsername": "admin", + "UptimeKumaPassword": "password", + } + + testCases := []struct { + name string + key string + value string + wantErr bool + }{ + {name: "enabled true", key: "UptimeKumaEnabled", value: "true"}, + {name: "enabled false", key: "UptimeKumaEnabled", value: "false"}, + {name: "enabled invalid", key: "UptimeKumaEnabled", value: "on", wantErr: true}, + {name: "url http valid", key: "UptimeKumaUrl", value: "http://192.168.1.100:3001"}, + {name: "url https valid", key: "UptimeKumaUrl", value: "https://kuma.example.com"}, + {name: "url invalid", key: "UptimeKumaUrl", value: "kuma.example.com", wantErr: true}, + {name: "scope all", key: "UptimeKumaMonitorScope", value: "all"}, + {name: "scope selected", key: "UptimeKumaMonitorScope", value: "selected"}, + {name: "scope invalid", key: "UptimeKumaMonitorScope", value: "none", wantErr: true}, + {name: "sync interval valid", key: "UptimeKumaSyncInterval", value: "5"}, + {name: "sync interval invalid", key: "UptimeKumaSyncInterval", value: "0", wantErr: true}, + {name: "interval valid", key: "UptimeKumaInterval", value: "60"}, + {name: "interval invalid", key: "UptimeKumaInterval", value: "-60", wantErr: true}, + {name: "retry valid", key: "UptimeKumaRetry", value: "0"}, + {name: "retry positive valid", key: "UptimeKumaRetry", value: "3"}, + {name: "retry invalid", key: "UptimeKumaRetry", value: "-1", wantErr: true}, + } + + for _, tc := range testCases { + err := validateUptimeKumaOption(tc.key, tc.value, state) + if tc.wantErr && err == nil { + t.Fatalf("%s: expected error", tc.name) + } + if !tc.wantErr && err != nil { + t.Fatalf("%s: unexpected error: %v", tc.name, err) + } + } + + // Test enabling Uptime Kuma when URL or credentials are empty in state + stateEmpty := map[string]string{ + "UptimeKumaUrl": "", + "UptimeKumaUsername": "", + "UptimeKumaPassword": "", + } + if err := validateUptimeKumaOption("UptimeKumaEnabled", "true", stateEmpty); err == nil { + t.Fatal("expected error when enabling Uptime Kuma with empty URL/credentials in state") + } +} diff --git a/openflare_server/controller/uptimekuma.go b/openflare_server/controller/uptimekuma.go new file mode 100644 index 00000000..c3b1cef4 --- /dev/null +++ b/openflare_server/controller/uptimekuma.go @@ -0,0 +1,24 @@ +package controller + +import ( + "openflare/service" + + "github.com/gin-gonic/gin" +) + +// SyncUptimeKuma godoc +// @Summary Manually trigger Uptime Kuma sync +// @Tags UptimeKuma +// @Accept json +// @Produce json +// @Security BearerAuth +// @Success 200 {object} map[string]interface{} +// @Router /api/uptimekuma/sync [post] +func SyncUptimeKuma(c *gin.Context) { + err := service.SyncToUptimeKuma() + if err != nil { + respondFailure(c, err.Error()) + return + } + respondSuccessMessage(c, "同步成功") +} diff --git a/openflare_server/job/cron.go b/openflare_server/job/cron.go index 2694c269..73416e12 100644 --- a/openflare_server/job/cron.go +++ b/openflare_server/job/cron.go @@ -26,6 +26,14 @@ func InitCronJobs() { slog.Info("registered WAF IP group sync cron job") } + // Register Uptime Kuma sync job (check every minute) + _, err = cronRunner.AddJob("* * * * *", &UptimeKumaSyncJob{}) + if err != nil { + slog.Error("failed to register Uptime Kuma sync cron job", "error", err) + } else { + slog.Info("registered Uptime Kuma sync cron job") + } + cronRunner.Start() } diff --git a/openflare_server/job/uptimekuma.go b/openflare_server/job/uptimekuma.go new file mode 100644 index 00000000..49578050 --- /dev/null +++ b/openflare_server/job/uptimekuma.go @@ -0,0 +1,43 @@ +package job + +import ( + "log/slog" + "openflare/common" + "openflare/service" + "sync" + "time" +) + +var lastUptimeKumaSyncTime time.Time +var uptimeKumaSyncMutex sync.Mutex + +type UptimeKumaSyncJob struct{} + +func (j *UptimeKumaSyncJob) Run() { + if !common.UptimeKumaEnabled { + return + } + + interval := common.UptimeKumaSyncInterval + if interval <= 0 { + interval = 5 + } + + if time.Since(lastUptimeKumaSyncTime) < time.Duration(interval)*time.Minute { + return + } + + if !uptimeKumaSyncMutex.TryLock() { + slog.Warn("Uptime Kuma sync job is already running, skipping this scheduled run") + return + } + defer uptimeKumaSyncMutex.Unlock() + + slog.Info("Starting scheduled Uptime Kuma sync") + if err := service.SyncToUptimeKuma(); err != nil { + slog.Error("Uptime Kuma sync failed", "error", err) + } else { + lastUptimeKumaSyncTime = time.Now() + slog.Info("Uptime Kuma sync completed successfully") + } +} diff --git a/openflare_server/model/option.go b/openflare_server/model/option.go index 3b96eca1..df188496 100644 --- a/openflare_server/model/option.go +++ b/openflare_server/model/option.go @@ -52,6 +52,17 @@ func InitOptionMap() { common.OptionMap["AgentUpdateRepo"] = common.AgentUpdateRepo common.OptionMap["GeoIPProvider"] = common.GeoIPProvider common.OptionMap["DatabaseAutoCleanupEnabled"] = strconv.FormatBool(common.DatabaseAutoCleanupEnabled) + common.OptionMap["UptimeKumaEnabled"] = strconv.FormatBool(common.UptimeKumaEnabled) + common.OptionMap["UptimeKumaUrl"] = common.UptimeKumaUrl + common.OptionMap["UptimeKumaUsername"] = common.UptimeKumaUsername + common.OptionMap["UptimeKumaPassword"] = common.UptimeKumaPassword + common.OptionMap["UptimeKumaMonitorScope"] = common.UptimeKumaMonitorScope + common.OptionMap["UptimeKumaSelectedSites"] = common.UptimeKumaSelectedSites + common.OptionMap["UptimeKumaSyncInterval"] = strconv.Itoa(common.UptimeKumaSyncInterval) + common.OptionMap["UptimeKumaInterval"] = strconv.Itoa(common.UptimeKumaInterval) + common.OptionMap["UptimeKumaRetry"] = strconv.Itoa(common.UptimeKumaRetry) + common.OptionMap["UptimeKumaRetryInterval"] = strconv.Itoa(common.UptimeKumaRetryInterval) + common.OptionMap["UptimeKumaTimeout"] = strconv.Itoa(common.UptimeKumaTimeout) common.OptionMap["DatabaseAutoCleanupRetentionDays"] = strconv.Itoa(common.DatabaseAutoCleanupRetentionDays) common.OptionMap["OpenRestyDefaultServerReturnStatus"] = strconv.Itoa(common.OpenRestyDefaultServerReturnStatus) common.OptionMap["OpenRestyWorkerProcesses"] = common.OpenRestyWorkerProcesses @@ -116,6 +127,9 @@ func UpdateOptions(options []Option) error { if err := DB.Transaction(func(tx *gorm.DB) error { for _, item := range options { + if item.Key == "UptimeKumaPassword" && strings.TrimSpace(item.Value) == "" { + continue + } option := Option{ Key: item.Key, } @@ -133,6 +147,9 @@ func UpdateOptions(options []Option) error { } for _, item := range options { + if item.Key == "UptimeKumaPassword" && strings.TrimSpace(item.Value) == "" { + continue + } updateOptionMap(item.Key, item.Value) } return nil @@ -209,6 +226,38 @@ func updateOptionMap(key string, value string) { common.GeoIPProvider = value shouldRefreshGeoIP = true } + case "UptimeKumaEnabled": + common.UptimeKumaEnabled = value == "true" + case "UptimeKumaUrl": + common.UptimeKumaUrl = value + case "UptimeKumaUsername": + common.UptimeKumaUsername = value + case "UptimeKumaPassword": + common.UptimeKumaPassword = value + case "UptimeKumaMonitorScope": + common.UptimeKumaMonitorScope = value + case "UptimeKumaSelectedSites": + common.UptimeKumaSelectedSites = value + case "UptimeKumaSyncInterval": + if v, err := strconv.Atoi(value); err == nil && v > 0 { + common.UptimeKumaSyncInterval = v + } + case "UptimeKumaInterval": + if v, err := strconv.Atoi(value); err == nil && v > 0 { + common.UptimeKumaInterval = v + } + case "UptimeKumaRetry": + if v, err := strconv.Atoi(value); err == nil && v >= 0 { + common.UptimeKumaRetry = v + } + case "UptimeKumaRetryInterval": + if v, err := strconv.Atoi(value); err == nil && v > 0 { + common.UptimeKumaRetryInterval = v + } + case "UptimeKumaTimeout": + if v, err := strconv.Atoi(value); err == nil && v > 0 { + common.UptimeKumaTimeout = v + } case "DatabaseAutoCleanupEnabled": common.DatabaseAutoCleanupEnabled = value == "true" case "DatabaseAutoCleanupRetentionDays": diff --git a/openflare_server/router/api-router.go b/openflare_server/router/api-router.go index fc8ad27e..3b247346 100644 --- a/openflare_server/router/api-router.go +++ b/openflare_server/router/api-router.go @@ -67,6 +67,11 @@ func SetApiRouter(router *gin.Engine) { optionRoute.POST("/geoip/lookup", controller.LookupGeoIP) optionRoute.POST("/database/cleanup", controller.CleanupDatabaseObservability) } + uptimekumaRoute := apiRouter.Group("/uptimekuma") + uptimekumaRoute.Use(middleware.RootAuth(), middleware.NoTokenAuth()) + { + uptimekumaRoute.POST("/sync", controller.SyncUptimeKuma) + } authSourceRoute := apiRouter.Group("/auth-sources") authSourceRoute.Use(middleware.RootAuth(), middleware.NoTokenAuth()) { diff --git a/openflare_server/router/api_uptimekuma_test.go b/openflare_server/router/api_uptimekuma_test.go new file mode 100644 index 00000000..4af20b64 --- /dev/null +++ b/openflare_server/router/api_uptimekuma_test.go @@ -0,0 +1,454 @@ +package router_test + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "openflare/common" + "openflare/model" + "openflare/router" + "strings" + "sync" + "testing" + "time" + + "github.com/gin-contrib/sessions" + "github.com/gin-contrib/sessions/cookie" + "github.com/gin-gonic/gin" +) + +// mockKumaServer simulates Uptime Kuma's Engine.IO/Socket.IO polling endpoints +type mockKumaServer struct { + mu sync.Mutex + postsReceived []string + pendingPackets chan string + monitorList string // JSON representing map[string]UptimeKumaMonitor +} + +func newMockKumaServer(monitorList string) *mockKumaServer { + return &mockKumaServer{ + pendingPackets: make(chan string, 100), + monitorList: monitorList, + } +} + +func (s *mockKumaServer) ServeHTTP(w http.ResponseWriter, r *http.Request) { + s.mu.Lock() + defer s.mu.Unlock() + + transport := r.URL.Query().Get("transport") + sid := r.URL.Query().Get("sid") + + if r.Method == "GET" { + if transport == "polling" && sid == "" { + // Handshake response + w.Header().Set("Content-Type", "text/plain;charset=UTF-8") + _, _ = w.Write([]byte(`0{"sid":"mock-sid"}`)) + return + } + + if transport == "polling" && sid == "mock-sid" { + // Long-polling GET request + w.Header().Set("Content-Type", "text/plain;charset=UTF-8") + select { + case pkt := <-s.pendingPackets: + _, _ = w.Write([]byte(pkt)) + case <-time.After(100 * time.Millisecond): + _, _ = w.Write([]byte("")) + } + return + } + } else if r.Method == "POST" { + bodyBytes, _ := io.ReadAll(r.Body) + bodyStr := string(bodyBytes) + s.postsReceived = append(s.postsReceived, bodyStr) + + w.Header().Set("Content-Type", "text/plain;charset=UTF-8") + w.WriteHeader(http.StatusOK) + + if bodyStr == "40" { + // Namespace Connect event + // Immediately queue the monitorList payload to be fetched by the next GET poll + s.pendingPackets <- fmt.Sprintf(`42["monitorList",%s]`, s.monitorList) + return + } + + if strings.HasPrefix(bodyStr, "42") { + // Socket.IO message: 42[...] + payload := bodyStr[2:] + // Find ack ID (digits at the start of payload) + digitsEnd := 0 + for digitsEnd < len(payload) && payload[digitsEnd] >= '0' && payload[digitsEnd] <= '9' { + digitsEnd++ + } + if digitsEnd == 0 { + return + } + ackIDStr := payload[:digitsEnd] + jsonArrayStr := payload[digitsEnd:] + + var arr []json.RawMessage + if err := json.Unmarshal([]byte(jsonArrayStr), &arr); err != nil || len(arr) == 0 { + return + } + + var eventName string + _ = json.Unmarshal(arr[0], &eventName) + + switch eventName { + case "login", "loginByToken": + s.pendingPackets <- fmt.Sprintf("43%s[{\"ok\":true}]", ackIDStr) + case "getTags": + s.pendingPackets <- fmt.Sprintf("43%s[{\"ok\":true,\"tags\":[{\"id\":10,\"name\":\"OpenFlare\",\"color\":\"#4f46e5\"}]}]", ackIDStr) + case "addTag": + s.pendingPackets <- fmt.Sprintf("43%s[{\"ok\":true,\"tag\":{\"id\":10}}]", ackIDStr) + case "add": + s.pendingPackets <- fmt.Sprintf("43%s[{\"ok\":true,\"monitorID\":100}]", ackIDStr) + case "addMonitorTag": + s.pendingPackets <- fmt.Sprintf("43%s[{\"ok\":true}]", ackIDStr) + case "editMonitor": + s.pendingPackets <- fmt.Sprintf("43%s[{\"ok\":true}]", ackIDStr) + case "deleteMonitor": + s.pendingPackets <- fmt.Sprintf("43%s[{\"ok\":true}]", ackIDStr) + } + } + } +} + +func TestUptimeKumaSyncDisabled(t *testing.T) { + gin.SetMode(gin.TestMode) + common.RedisEnabled = false + setupTestDB(t) + + engine := gin.New() + engine.Use(sessions.Sessions("session", cookie.NewStore([]byte("test-secret")))) + router.SetApiRouter(engine) + + loginCookie := loginAsRoot(t, engine) + + // Keep integration disabled + common.UptimeKumaEnabled = false + + // Request sync, should fail + req := httptest.NewRequest(http.MethodPost, "/api/uptimekuma/sync", nil) + req.AddCookie(loginCookie) + recorder := httptest.NewRecorder() + engine.ServeHTTP(recorder, req) + + if recorder.Code != http.StatusOK { + t.Fatalf("expected status 200, got %d", recorder.Code) + } + + var resp apiResponse + if err := json.Unmarshal(recorder.Body.Bytes(), &resp); err != nil { + t.Fatalf("failed to decode response: %v", err) + } + + if resp.Success { + t.Fatal("expected sync request to fail when integration is disabled") + } + if !strings.Contains(resp.Message, "disabled") { + t.Fatalf("expected error message to mention integration is disabled, got: %s", resp.Message) + } +} + +func TestUptimeKumaSyncSuccess(t *testing.T) { + gin.SetMode(gin.TestMode) + common.RedisEnabled = false + setupTestDB(t) + + // Clean up route table just in case + _ = model.DB.Where("1 = 1").Delete(&model.ProxyRoute{}).Error + + // Seed proxy routes + // Route 1: site-a (exists in Uptime Kuma but has different check parameters - should trigger editMonitor) + routeA := &model.ProxyRoute{ + SiteName: "site-a", + Domain: "site-a.com", + Domains: `["site-a.com"]`, + OriginURL: "http://10.0.0.1", + Enabled: true, + EnableHTTPS: false, + } + // Route 2: site-b (does not exist in Uptime Kuma - should trigger add & addMonitorTag) + routeB := &model.ProxyRoute{ + SiteName: "site-b", + Domain: "site-b.com", + Domains: `["site-b.com"]`, + OriginURL: "https://10.0.0.2", + Enabled: true, + EnableHTTPS: true, + } + // Route 3: site-c (disabled locally - should NOT be processed/created) + routeC := &model.ProxyRoute{ + SiteName: "site-c", + Domain: "site-c.com", + Domains: `["site-c.com"]`, + OriginURL: "http://10.0.0.3", + Enabled: false, + EnableHTTPS: false, + } + + if err := model.DB.Create(routeA).Error; err != nil { + t.Fatalf("failed to seed routeA: %v", err) + } + if err := model.DB.Create(routeB).Error; err != nil { + t.Fatalf("failed to seed routeB: %v", err) + } + if err := model.DB.Create(routeC).Error; err != nil { + t.Fatalf("failed to seed routeC: %v", err) + } + + // Prepare mock monitorList + // 1. "site-old": tagged with OpenFlare but doesn't exist locally anymore -> should trigger deleteMonitor + // 2. "site-a": matches routeA but has interval = 30 (default UptimeKumaInterval is 60) -> should trigger editMonitor + monitorListJSON := `{ + "99": { + "id": 99, + "name": "site-old", + "url": "http://site-old.com", + "interval": 60, + "tags": [{"tag_id": 10, "name": "OpenFlare"}] + }, + "98": { + "id": 98, + "name": "site-a", + "url": "http://site-a.com", + "interval": 30, + "tags": [{"tag_id": 10, "name": "OpenFlare"}] + } + }` + + mockSrv := newMockKumaServer(monitorListJSON) + server := httptest.NewServer(mockSrv) + defer server.Close() + + // Backup and set configs + oldEnabled := common.UptimeKumaEnabled + oldUrl := common.UptimeKumaUrl + oldUsername := common.UptimeKumaUsername + oldPassword := common.UptimeKumaPassword + oldScope := common.UptimeKumaMonitorScope + oldInterval := common.UptimeKumaInterval + oldRetry := common.UptimeKumaRetry + oldRetryInterval := common.UptimeKumaRetryInterval + oldTimeout := common.UptimeKumaTimeout + + common.UptimeKumaEnabled = true + common.UptimeKumaUrl = server.URL + common.UptimeKumaUsername = "admin" + common.UptimeKumaPassword = "password" + common.UptimeKumaMonitorScope = "all" + common.UptimeKumaInterval = 60 + common.UptimeKumaRetry = 0 + common.UptimeKumaRetryInterval = 60 + common.UptimeKumaTimeout = 48 + + defer func() { + common.UptimeKumaEnabled = oldEnabled + common.UptimeKumaUrl = oldUrl + common.UptimeKumaUsername = oldUsername + common.UptimeKumaPassword = oldPassword + common.UptimeKumaMonitorScope = oldScope + common.UptimeKumaInterval = oldInterval + common.UptimeKumaRetry = oldRetry + common.UptimeKumaRetryInterval = oldRetryInterval + common.UptimeKumaTimeout = oldTimeout + }() + + engine := gin.New() + engine.Use(sessions.Sessions("session", cookie.NewStore([]byte("test-secret")))) + router.SetApiRouter(engine) + + loginCookie := loginAsRoot(t, engine) + + req := httptest.NewRequest(http.MethodPost, "/api/uptimekuma/sync", nil) + req.AddCookie(loginCookie) + recorder := httptest.NewRecorder() + engine.ServeHTTP(recorder, req) + + if recorder.Code != http.StatusOK { + t.Fatalf("expected status 200, got %d. Body: %s", recorder.Code, recorder.Body.String()) + } + + var resp apiResponse + if err := json.Unmarshal(recorder.Body.Bytes(), &resp); err != nil { + t.Fatalf("failed to decode response: %v", err) + } + + if !resp.Success { + t.Fatalf("sync request failed: %s", resp.Message) + } + + mockSrv.mu.Lock() + posts := mockSrv.postsReceived + mockSrv.mu.Unlock() + + // Verify events received + hasLogin := false + hasGetTags := false + hasAddSiteB := false + hasTagSiteB := false + hasEditSiteA := false + hasDeleteOld := false + + for _, body := range posts { + if strings.Contains(body, `"login"`) && strings.Contains(body, `"admin"`) && strings.Contains(body, `"password"`) { + hasLogin = true + } + if strings.Contains(body, `"getTags"`) { + hasGetTags = true + } + if strings.Contains(body, `"add"`) && strings.Contains(body, `"site-b"`) && strings.Contains(body, `"https://site-b.com"`) { + hasAddSiteB = true + } + if strings.Contains(body, `"addMonitorTag"`) && strings.Contains(body, `10`) && strings.Contains(body, `100`) { + hasTagSiteB = true + } + if strings.Contains(body, `"editMonitor"`) && strings.Contains(body, `98`) && strings.Contains(body, `"site-a"`) && strings.Contains(body, `"interval":60`) { + hasEditSiteA = true + } + if strings.Contains(body, `"deleteMonitor"`) && strings.Contains(body, `99`) { + hasDeleteOld = true + } + } + + if !hasLogin { + t.Error("expected login event to be called") + } + if !hasGetTags { + t.Error("expected getTags event to be called") + } + if !hasAddSiteB { + t.Error("expected site-b to be added") + } + if !hasTagSiteB { + t.Error("expected site-b to be tagged") + } + if !hasEditSiteA { + t.Error("expected site-a to be edited/updated") + } + if !hasDeleteOld { + t.Error("expected site-old to be deleted") + } +} + +func TestUptimeKumaSyncSelectedScope(t *testing.T) { + gin.SetMode(gin.TestMode) + common.RedisEnabled = false + setupTestDB(t) + + // Clean up route table + _ = model.DB.Where("1 = 1").Delete(&model.ProxyRoute{}).Error + + // Seed proxy routes + // Route 1: site-a (enabled, in selected list) + routeA := &model.ProxyRoute{ + SiteName: "site-a", + Domain: "site-a.com", + Domains: `["site-a.com"]`, + OriginURL: "http://10.0.0.1", + Enabled: true, + EnableHTTPS: false, + } + // Route 2: site-b (enabled, NOT in selected list) + routeB := &model.ProxyRoute{ + SiteName: "site-b", + Domain: "site-b.com", + Domains: `["site-b.com"]`, + OriginURL: "http://10.0.0.2", + Enabled: true, + EnableHTTPS: false, + } + + if err := model.DB.Create(routeA).Error; err != nil { + t.Fatalf("failed to seed routeA: %v", err) + } + if err := model.DB.Create(routeB).Error; err != nil { + t.Fatalf("failed to seed routeB: %v", err) + } + + mockSrv := newMockKumaServer(`{}`) + server := httptest.NewServer(mockSrv) + defer server.Close() + + // Backup and set configs + oldEnabled := common.UptimeKumaEnabled + oldUrl := common.UptimeKumaUrl + oldUsername := common.UptimeKumaUsername + oldPassword := common.UptimeKumaPassword + oldScope := common.UptimeKumaMonitorScope + oldSelected := common.UptimeKumaSelectedSites + + common.UptimeKumaEnabled = true + common.UptimeKumaUrl = server.URL + common.UptimeKumaUsername = "admin" + common.UptimeKumaPassword = "password" + common.UptimeKumaMonitorScope = "selected" + common.UptimeKumaSelectedSites = "site-a" // site-b is excluded + + defer func() { + common.UptimeKumaEnabled = oldEnabled + common.UptimeKumaUrl = oldUrl + common.UptimeKumaUsername = oldUsername + common.UptimeKumaPassword = oldPassword + common.UptimeKumaMonitorScope = oldScope + common.UptimeKumaSelectedSites = oldSelected + }() + + engine := gin.New() + engine.Use(sessions.Sessions("session", cookie.NewStore([]byte("test-secret")))) + router.SetApiRouter(engine) + + loginCookie := loginAsRoot(t, engine) + + req := httptest.NewRequest(http.MethodPost, "/api/uptimekuma/sync", nil) + req.AddCookie(loginCookie) + recorder := httptest.NewRecorder() + engine.ServeHTTP(recorder, req) + + if recorder.Code != http.StatusOK { + t.Fatalf("expected status 200, got %d", recorder.Code) + } + + var resp apiResponse + if err := json.Unmarshal(recorder.Body.Bytes(), &resp); err != nil { + t.Fatalf("failed to decode response: %v", err) + } + + if !resp.Success { + t.Fatalf("sync request failed: %s", resp.Message) + } + + mockSrv.mu.Lock() + posts := mockSrv.postsReceived + mockSrv.mu.Unlock() + + hasLogin := false + hasAddSiteA := false + hasAddSiteB := false + + for _, body := range posts { + if strings.Contains(body, `"login"`) && strings.Contains(body, `"admin"`) && strings.Contains(body, `"password"`) { + hasLogin = true + } + if strings.Contains(body, `"add"`) && strings.Contains(body, `"site-a"`) { + hasAddSiteA = true + } + if strings.Contains(body, `"add"`) && strings.Contains(body, `"site-b"`) { + hasAddSiteB = true + } + } + + if !hasLogin { + t.Error("expected login event to be called") + } + if !hasAddSiteA { + t.Error("expected site-a to be added") + } + if hasAddSiteB { + t.Error("expected site-b NOT to be added (not in selected scope)") + } +} diff --git a/openflare_server/service/uptimekuma.go b/openflare_server/service/uptimekuma.go new file mode 100644 index 00000000..04e617ee --- /dev/null +++ b/openflare_server/service/uptimekuma.go @@ -0,0 +1,302 @@ +package service + +import ( + "fmt" + "log/slog" + "openflare/common" + "openflare/model" + "openflare/utils/uptimekuma" + "strings" + "sync/atomic" + "time" +) + +var isSyncing atomic.Bool + +func SyncToUptimeKuma() error { + if !common.UptimeKumaEnabled { + return fmt.Errorf("Uptime Kuma integration is disabled") + } + + if !isSyncing.CompareAndSwap(false, true) { + return fmt.Errorf("sync task is already in progress, please try again later") + } + defer isSyncing.Store(false) + + kumaUrl := strings.TrimSpace(common.UptimeKumaUrl) + kumaUsername := strings.TrimSpace(common.UptimeKumaUsername) + kumaPassword := strings.TrimSpace(common.UptimeKumaPassword) + if kumaUrl == "" || kumaUsername == "" || kumaPassword == "" { + return fmt.Errorf("Uptime Kuma URL, username, or password is not configured (URL: %q, Username: %q, PasswordLength: %d)", kumaUrl, kumaUsername, len(kumaPassword)) + } + + slog.Info("Starting Uptime Kuma sync process", "url", kumaUrl, "username", kumaUsername, "scope", common.UptimeKumaMonitorScope) + + // 1. Fetch expected sites + allRoutes, err := model.ListProxyRoutes() + if err != nil { + return fmt.Errorf("failed to list local proxy routes: %w", err) + } + + var expectedRoutes []*model.ProxyRoute + scope := common.UptimeKumaMonitorScope + if scope == "selected" { + selectedList := strings.Split(common.UptimeKumaSelectedSites, ",") + selectedMap := make(map[string]bool) + for _, name := range selectedList { + trimmedName := strings.TrimSpace(name) + if trimmedName != "" { + selectedMap[trimmedName] = true + } + } + for _, route := range allRoutes { + if route.Enabled && selectedMap[route.SiteName] { + expectedRoutes = append(expectedRoutes, route) + } + } + } else { + for _, route := range allRoutes { + if route.Enabled { + expectedRoutes = append(expectedRoutes, route) + } + } + } + + // 2. Connect to Uptime Kuma + slog.Debug("Connecting to Uptime Kuma socket endpoint", "url", kumaUrl) + client := uptimekuma.NewSocketIOClient(kumaUrl) + if err := client.Connect(); err != nil { + slog.Error("Failed to connect to Uptime Kuma endpoint", "url", kumaUrl, "error", err) + return fmt.Errorf("failed to connect to Uptime Kuma: %w", err) + } + defer client.Close() + + // 3. Login + slog.Debug("Sending login request to Uptime Kuma", "username", kumaUsername) + var loginAck string + loginPayload := map[string]string{ + "username": kumaUsername, + "password": kumaPassword, + } + loginAck, err = client.Emit("login", loginPayload) + if err != nil { + slog.Error("Failed to send login request to Uptime Kuma", "username", kumaUsername, "error", err) + return fmt.Errorf("login request failed: %w", err) + } + + var loginResult struct { + Ok bool `json:"ok"` + } + if err := uptimekuma.ParseAckResponse(loginAck, &loginResult); err != nil || !loginResult.Ok { + slog.Error("Uptime Kuma login verification failed", "username", kumaUsername, "error", err) + return fmt.Errorf("login failed: %w", err) + } + slog.Debug("Successfully logged into Uptime Kuma", "username", kumaUsername) + + // 4. Wait for monitor list event + slog.Debug("Waiting for monitor list push from Uptime Kuma") + select { + case <-client.GetMonitorListChan(): + slog.Debug("Received monitor list from Uptime Kuma") + case <-time.After(5 * time.Second): + slog.Error("Timeout waiting for Uptime Kuma monitorList push event") + return fmt.Errorf("timeout waiting for monitorList event from Uptime Kuma") + } + + // 5. Get existing tags to find "OpenFlare" + slog.Debug("Fetching tags from Uptime Kuma") + tagsAck, err := client.Emit("getTags") + if err != nil { + slog.Error("Failed to request tags from Uptime Kuma", "error", err) + return fmt.Errorf("failed to fetch tags: %w", err) + } + + var tagsResult struct { + Ok bool `json:"ok"` + Tags []uptimekuma.UptimeKumaTagItem `json:"tags"` + } + if err := uptimekuma.ParseAckResponse(tagsAck, &tagsResult); err != nil { + slog.Error("Failed to parse tags response from Uptime Kuma", "error", err) + return fmt.Errorf("parse tags response failed: %w", err) + } + + var openFlareTagID int + for _, t := range tagsResult.Tags { + if t.Name == "OpenFlare" { + openFlareTagID = t.ID + break + } + } + + // Create "OpenFlare" tag if not exists + if openFlareTagID == 0 { + slog.Debug("OpenFlare tag not found, creating new tag") + addTagAck, err := client.Emit("addTag", map[string]string{ + "name": "OpenFlare", + "color": "#4f46e5", + }) + if err != nil { + slog.Error("Failed to create OpenFlare tag in Uptime Kuma", "error", err) + return fmt.Errorf("failed to create tag: %w", err) + } + var tagResult struct { + Ok bool `json:"ok"` + Tag struct { + ID int `json:"id"` + } `json:"tag"` + } + if err := uptimekuma.ParseAckResponse(addTagAck, &tagResult); err != nil || tagResult.Tag.ID == 0 { + slog.Error("Failed to parse addTag response from Uptime Kuma", "error", err) + return fmt.Errorf("parse addTag response failed: %w", err) + } + openFlareTagID = tagResult.Tag.ID + slog.Debug("Successfully created OpenFlare tag", "tag_id", openFlareTagID) + } else { + slog.Debug("Found existing OpenFlare tag", "tag_id", openFlareTagID) + } + + // 6. Filter existing monitors by "OpenFlare" tag + existingOpenFlareMonitors := make(map[string]uptimekuma.UptimeKumaMonitor) + monitors := client.GetMonitorList() + for _, m := range monitors { + hasOpenFlareTag := false + for _, tag := range m.Tags { + if tag.Name == "OpenFlare" || tag.ID == openFlareTagID { + hasOpenFlareTag = true + break + } + } + if hasOpenFlareTag { + existingOpenFlareMonitors[m.Name] = m + } + } + + // Helper to format route URL + getRouteURL := func(route *model.ProxyRoute) string { + domains, err := decodeStoredDomains(route.Domains, route.Domain) + domain := route.Domain + if err == nil && len(domains) > 0 { + domain = domains[0] + } + if route.EnableHTTPS { + return "https://" + domain + } + return "http://" + domain + } + + expectedSitesMap := make(map[string]bool) + + // 7. Sync Loop + for _, route := range expectedRoutes { + expectedSitesMap[route.SiteName] = true + targetURL := getRouteURL(route) + + existing, exists := existingOpenFlareMonitors[route.SiteName] + if !exists { + // Create monitor + slog.Info("Creating monitor in Uptime Kuma", "name", route.SiteName, "url", targetURL) + monitorPayload := map[string]any{ + "type": "http", + "name": route.SiteName, + "url": targetURL, + "interval": common.UptimeKumaInterval, + "maxretries": common.UptimeKumaRetry, + "retryInterval": common.UptimeKumaRetryInterval, + "timeout": common.UptimeKumaTimeout, + "active": true, + "resendInterval": 0, + "expiryNotification": false, + "ignoreTls": false, + "accepted_statuscodes": []string{"200-299"}, + "dns_resolve_type": "A", + } + addAck, err := client.Emit("add", monitorPayload) + if err != nil { + slog.Error("Failed to add monitor to Uptime Kuma", "name", route.SiteName, "error", err) + continue + } + var addResult struct { + Ok bool `json:"ok"` + MonitorID int `json:"monitorID"` + } + if err := uptimekuma.ParseAckResponse(addAck, &addResult); err != nil || addResult.MonitorID == 0 { + slog.Error("Failed to parse add monitor result", "name", route.SiteName, "error", err) + continue + } + + // Add tag + slog.Debug("Adding OpenFlare tag to the new monitor", "name", route.SiteName, "monitor_id", addResult.MonitorID, "tag_id", openFlareTagID) + tagAck, err := client.Emit("addMonitorTag", openFlareTagID, addResult.MonitorID, "") + if err != nil { + slog.Error("Failed to add tag to monitor in Uptime Kuma", "name", route.SiteName, "monitorID", addResult.MonitorID, "error", err) + } else { + if err := uptimekuma.ParseAckResponse(tagAck, nil); err != nil { + slog.Error("Failed to parse add tag result", "name", route.SiteName, "monitorID", addResult.MonitorID, "error", err) + } else { + slog.Debug("OpenFlare tag successfully added to monitor", "name", route.SiteName, "monitor_id", addResult.MonitorID) + } + } + } else { + // Check if updates are needed + needsUpdate := existing.Url != targetURL || + existing.Interval != common.UptimeKumaInterval || + existing.MaxRetries != common.UptimeKumaRetry || + existing.RetryInterval != common.UptimeKumaRetryInterval || + existing.Timeout != common.UptimeKumaTimeout + + if needsUpdate { + slog.Info("Updating monitor in Uptime Kuma due to settings mismatch", + "name", route.SiteName, + "url_changed", existing.Url != targetURL, + "interval_changed", existing.Interval != common.UptimeKumaInterval, + "max_retries_changed", existing.MaxRetries != common.UptimeKumaRetry, + "retry_interval_changed", existing.RetryInterval != common.UptimeKumaRetryInterval, + "timeout_changed", existing.Timeout != common.UptimeKumaTimeout, + ) + monitorPayload := map[string]any{ + "id": existing.ID, + "type": "http", + "name": route.SiteName, + "url": targetURL, + "interval": common.UptimeKumaInterval, + "maxretries": common.UptimeKumaRetry, + "retryInterval": common.UptimeKumaRetryInterval, + "timeout": common.UptimeKumaTimeout, + "active": true, + "resendInterval": 0, + "expiryNotification": false, + "ignoreTls": false, + "accepted_statuscodes": []string{"200-299"}, + "dns_resolve_type": "A", + } + editAck, err := client.Emit("editMonitor", monitorPayload) + if err != nil { + slog.Error("Failed to edit monitor in Uptime Kuma", "name", route.SiteName, "error", err) + } else { + if err := uptimekuma.ParseAckResponse(editAck, nil); err != nil { + slog.Error("Failed to parse edit monitor result", "name", route.SiteName, "error", err) + } else { + slog.Info("Successfully updated monitor in Uptime Kuma", "name", route.SiteName) + } + } + } + } + } + + // 8. Delete Loop + for name, m := range existingOpenFlareMonitors { + if !expectedSitesMap[name] { + slog.Info("Deleting monitor in Uptime Kuma", "name", name, "monitorID", m.ID) + deleteAck, err := client.Emit("deleteMonitor", m.ID) + if err != nil { + slog.Error("Failed to delete monitor in Uptime Kuma", "name", name, "monitorID", m.ID, "error", err) + } else { + if err := uptimekuma.ParseAckResponse(deleteAck, nil); err != nil { + slog.Error("Failed to parse delete monitor result", "name", name, "monitorID", m.ID, "error", err) + } + } + } + } + + return nil +} diff --git a/openflare_server/utils/uptimekuma/client.go b/openflare_server/utils/uptimekuma/client.go new file mode 100644 index 00000000..b591f3ad --- /dev/null +++ b/openflare_server/utils/uptimekuma/client.go @@ -0,0 +1,373 @@ +package uptimekuma + +import ( + "context" + "encoding/json" + "fmt" + "io" + "log/slog" + "net/http" + "strconv" + "strings" + "sync" + "time" +) + +type UptimeKumaMonitor struct { + ID int `json:"id"` + Name string `json:"name"` + Url string `json:"url"` + Type string `json:"type"` + Interval int `json:"interval"` + MaxRetries int `json:"maxretries"` + RetryInterval int `json:"retryInterval"` + Timeout int `json:"timeout"` + Tags []UptimeKumaTag `json:"tags"` +} + +type UptimeKumaTag struct { + ID int `json:"tag_id"` + Name string `json:"name"` + Color string `json:"color"` +} + +type UptimeKumaTagItem struct { + ID int `json:"id"` + Name string `json:"name"` + Color string `json:"color"` +} + +type SocketIOClient struct { + baseURL string + httpClient *http.Client + sid string + ackMutex sync.Mutex + ackID int + ackChanMap map[int]chan string + doneChan chan struct{} + closeOnce sync.Once + + monitorListMutex sync.RWMutex + monitorList map[string]UptimeKumaMonitor + monitorListChan chan struct{} + monitorListOnce sync.Once + + ctx context.Context + cancel context.CancelFunc + + err error +} + +func NewSocketIOClient(baseURL string) *SocketIOClient { + ctx, cancel := context.WithCancel(context.Background()) + return &SocketIOClient{ + baseURL: strings.TrimSuffix(baseURL, "/"), + httpClient: &http.Client{ + Timeout: 60 * time.Second, + }, + ackChanMap: make(map[int]chan string), + doneChan: make(chan struct{}), + monitorListChan: make(chan struct{}), + monitorList: make(map[string]UptimeKumaMonitor), + ctx: ctx, + cancel: cancel, + } +} + +func (c *SocketIOClient) Connect() error { + slog.Debug("Uptime Kuma client starting handshake", "baseURL", c.baseURL) + // 1. Handshake + u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling", c.baseURL) + reqHandshake, err := http.NewRequestWithContext(c.ctx, "GET", u, nil) + if err != nil { + return fmt.Errorf("create handshake request failed: %w", err) + } + resp, err := c.httpClient.Do(reqHandshake) + if err != nil { + slog.Error("Uptime Kuma handshake connection failed", "url", u, "error", err) + return fmt.Errorf("handshake request failed: %w", err) + } + defer resp.Body.Close() + + bs, err := io.ReadAll(resp.Body) + if err != nil { + slog.Error("Failed to read Uptime Kuma handshake response body", "error", err) + return fmt.Errorf("read handshake body failed: %w", err) + } + + bodyStr := string(bs) + slog.Debug("Received handshake response from Uptime Kuma", "body", bodyStr) + if len(bodyStr) == 0 || bodyStr[0] != '0' { + return fmt.Errorf("invalid handshake response format: %s", bodyStr) + } + + var hs struct { + Sid string `json:"sid"` + } + if err := json.Unmarshal([]byte(bodyStr[1:]), &hs); err != nil { + return fmt.Errorf("unmarshal handshake sid failed: %w", err) + } + c.sid = hs.Sid + slog.Debug("Uptime Kuma handshake success", "sid", c.sid) + + // 2. Namespace Connect + slog.Debug("Sending namespace connect request to Uptime Kuma", "sid", c.sid) + connectURL := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid) + req, err := http.NewRequestWithContext(c.ctx, "POST", connectURL, strings.NewReader("40")) + if err != nil { + return fmt.Errorf("create connect request failed: %w", err) + } + req.Header.Set("Content-Type", "text/plain;charset=UTF-8") + respConnect, err := c.httpClient.Do(req) + if err != nil { + slog.Error("Uptime Kuma namespace connect request failed", "sid", c.sid, "error", err) + return fmt.Errorf("namespace connect failed: %w", err) + } + respConnect.Body.Close() + slog.Debug("Namespace connected successfully to Uptime Kuma", "sid", c.sid) + + // 3. Start Polling Loop + go c.pollLoop() + + return nil +} + +func (c *SocketIOClient) pollLoop() { + slog.Debug("Uptime Kuma polling loop started", "sid", c.sid) + defer c.Close() + for { + select { + case <-c.doneChan: + slog.Debug("Uptime Kuma polling loop stopped (doneChan closed)", "sid", c.sid) + return + default: + } + + u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid) + reqPoll, err := http.NewRequestWithContext(c.ctx, "GET", u, nil) + if err != nil { + slog.Error("Failed to create Uptime Kuma polling request", "sid", c.sid, "error", err) + c.err = err + return + } + resp, err := c.httpClient.Do(reqPoll) + if err != nil { + slog.Error("Uptime Kuma polling request failed", "sid", c.sid, "error", err) + c.err = err + return + } + + bs, err := io.ReadAll(resp.Body) + resp.Body.Close() + if err != nil { + slog.Error("Failed to read Uptime Kuma polling body", "sid", c.sid, "error", err) + c.err = err + return + } + + bodyStr := string(bs) + if len(bodyStr) == 0 { + continue + } + + slog.Debug("Received polling payload from Uptime Kuma", "length", len(bodyStr)) + packets := strings.Split(bodyStr, "\x1e") + for _, pkt := range packets { + if len(pkt) == 0 { + continue + } + engineIOType := pkt[0] + payload := pkt[1:] + + slog.Debug("Parsing engine.io packet", "type", string(engineIOType), "payload_len", len(payload)) + switch engineIOType { + case '2': // Ping + slog.Debug("Received engine.io ping, responding with pong", "sid", c.sid) + c.sendPong() + case '4': // Message + if len(payload) == 0 { + continue + } + socketIOType := payload[0] + socketIOPayload := payload[1:] + + slog.Debug("Parsing socket.io packet", "type", string(socketIOType), "payload", socketIOPayload) + switch socketIOType { + case '2': // Event + c.handleEvent(socketIOPayload) + case '3': // Ack + c.handleAck(socketIOPayload) + } + } + } + } +} + +func (c *SocketIOClient) sendPong() { + u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid) + req, err := http.NewRequestWithContext(c.ctx, "POST", u, strings.NewReader("3")) + if err != nil { + return + } + req.Header.Set("Content-Type", "text/plain;charset=UTF-8") + resp, err := c.httpClient.Do(req) + if err == nil { + resp.Body.Close() + } +} + +func (c *SocketIOClient) handleEvent(payload string) { + var arr []json.RawMessage + if err := json.Unmarshal([]byte(payload), &arr); err != nil || len(arr) < 2 { + return + } + var eventName string + if err := json.Unmarshal(arr[0], &eventName); err != nil { + return + } + if eventName == "monitorList" { + var list map[string]UptimeKumaMonitor + if err := json.Unmarshal(arr[1], &list); err == nil { + c.monitorListMutex.Lock() + c.monitorList = list + c.monitorListMutex.Unlock() + c.monitorListOnce.Do(func() { + close(c.monitorListChan) + }) + } + } +} + +func (c *SocketIOClient) handleAck(payload string) { + idx := strings.IndexByte(payload, '[') + if idx == -1 { + return + } + ackIDStr := payload[:idx] + ackID, err := strconv.Atoi(ackIDStr) + if err != nil { + return + } + c.ackMutex.Lock() + ch, ok := c.ackChanMap[ackID] + if ok { + delete(c.ackChanMap, ackID) + c.ackMutex.Unlock() + select { + case ch <- payload[idx:]: + default: + } + } else { + c.ackMutex.Unlock() + } +} + +func (c *SocketIOClient) Emit(event string, args ...any) (string, error) { + c.ackMutex.Lock() + id := c.ackID + c.ackID++ + ch := make(chan string, 1) + c.ackChanMap[id] = ch + c.ackMutex.Unlock() + + payloadArr := []any{event} + payloadArr = append(payloadArr, args...) + bs, err := json.Marshal(payloadArr) + if err != nil { + c.ackMutex.Lock() + delete(c.ackChanMap, id) + c.ackMutex.Unlock() + slog.Error("Failed to marshal event payload", "event", event, "error", err) + return "", err + } + + body := fmt.Sprintf("42%d%s", id, string(bs)) + slog.Debug("Emitting Socket.IO event", "event", event, "ackID", id, "payload", string(bs)) + + u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid) + req, err := http.NewRequestWithContext(c.ctx, "POST", u, strings.NewReader(body)) + if err != nil { + c.ackMutex.Lock() + delete(c.ackChanMap, id) + c.ackMutex.Unlock() + return "", err + } + req.Header.Set("Content-Type", "text/plain;charset=UTF-8") + + resp, err := c.httpClient.Do(req) + if err != nil { + c.ackMutex.Lock() + delete(c.ackChanMap, id) + c.ackMutex.Unlock() + slog.Error("Failed to send Emit request", "event", event, "ackID", id, "error", err) + return "", err + } + resp.Body.Close() + + select { + case result := <-ch: + slog.Debug("Received Ack for event", "event", event, "ackID", id, "response", result) + return result, nil + case <-time.After(10 * time.Second): + c.ackMutex.Lock() + delete(c.ackChanMap, id) + c.ackMutex.Unlock() + slog.Error("Timeout waiting for event Ack", "event", event, "ackID", id) + return "", fmt.Errorf("timeout waiting for ack for event: %s", event) + case <-c.doneChan: + c.ackMutex.Lock() + delete(c.ackChanMap, id) + c.ackMutex.Unlock() + slog.Error("Client closed while waiting for event Ack", "event", event, "ackID", id) + return "", fmt.Errorf("client closed while waiting for event ack: %s", event) + } +} + +func (c *SocketIOClient) Close() { + c.closeOnce.Do(func() { + c.cancel() + close(c.doneChan) + }) +} + +func (c *SocketIOClient) GetMonitorListChan() <-chan struct{} { + return c.monitorListChan +} + +func (c *SocketIOClient) GetMonitorList() map[string]UptimeKumaMonitor { + c.monitorListMutex.RLock() + defer c.monitorListMutex.RUnlock() + + // Return a copy to prevent concurrent map read/write access + m := make(map[string]UptimeKumaMonitor, len(c.monitorList)) + for k, v := range c.monitorList { + m[k] = v + } + return m +} + +func ParseAckResponse(response string, target any) error { + var arr []json.RawMessage + if err := json.Unmarshal([]byte(response), &arr); err != nil || len(arr) == 0 { + return fmt.Errorf("invalid ack response format: %s", response) + } + + var status struct { + Ok bool `json:"ok"` + Msg string `json:"msg"` + } + if err := json.Unmarshal(arr[0], &status); err == nil { + if !status.Ok { + errMsg := status.Msg + if errMsg == "" { + errMsg = "unknown error from Uptime Kuma" + } + return fmt.Errorf("Uptime Kuma error response: %s", errMsg) + } + } + + if target != nil { + return json.Unmarshal(arr[0], target) + } + return nil +} diff --git a/openflare_server/web/features/settings/api/settings.ts b/openflare_server/web/features/settings/api/settings.ts index 15257a8f..d6a0492f 100644 --- a/openflare_server/web/features/settings/api/settings.ts +++ b/openflare_server/web/features/settings/api/settings.ts @@ -126,3 +126,9 @@ export function bindEmail(email: string, code: string) { export function getAboutContent() { return apiRequest('/about'); } + +export function syncUptimeKuma() { + return apiRequest('/uptimekuma/sync', { + method: 'POST', + }); +} diff --git a/openflare_server/web/features/settings/components/settings-page.tsx b/openflare_server/web/features/settings/components/settings-page.tsx index d8b4984c..a618f3db 100644 --- a/openflare_server/web/features/settings/components/settings-page.tsx +++ b/openflare_server/web/features/settings/components/settings-page.tsx @@ -31,8 +31,10 @@ import { rotateBootstrapToken, updateOptions, updateSelf, + syncUptimeKuma, } from '@/features/settings/api/settings'; import { AuthSourceModal } from '@/features/settings/components/auth-source-modal'; +import { UptimeKumaSiteSelectModal } from './uptimekuma-modal'; import type { BootstrapTokenPayload, DatabaseCleanupResult, @@ -87,6 +89,17 @@ const defaultOperationFields = { NodeOfflineThreshold: '120000', AgentUpdateRepo: 'Rain-kl/OpenFlare', GeoIPProvider: 'ipinfo', + UptimeKumaEnabled: false, + UptimeKumaUrl: '', + UptimeKumaUsername: '', + UptimeKumaPassword: '', + UptimeKumaMonitorScope: 'all', + UptimeKumaSelectedSites: '', + UptimeKumaSyncInterval: '5', + UptimeKumaInterval: '60', + UptimeKumaRetry: '0', + UptimeKumaRetryInterval: '60', + UptimeKumaTimeout: '48', OpenRestyDefaultServerReturnStatus: '421', OpenRestyWorkerProcesses: 'auto', OpenRestyWorkerConnections: '4096', @@ -255,6 +268,7 @@ export function SettingsPage() { const [cleanupModalState, setCleanupModalState] = useState(null); const [cleanupRetentionDays, setCleanupRetentionDays] = useState(''); + const [uptimeKumaModalOpen, setUptimeKumaModalOpen] = useState(false); const isRoot = (user?.role ?? 0) >= 100; @@ -428,6 +442,17 @@ export function SettingsPage() { GlobalWebRateLimitDuration: optionMap.GlobalWebRateLimitDuration ?? '180', CriticalRateLimitNum: optionMap.CriticalRateLimitNum ?? '100', CriticalRateLimitDuration: optionMap.CriticalRateLimitDuration ?? '1200', + UptimeKumaEnabled: toBoolean(optionMap.UptimeKumaEnabled, false), + UptimeKumaUrl: optionMap.UptimeKumaUrl ?? '', + UptimeKumaUsername: optionMap.UptimeKumaUsername ?? '', + UptimeKumaPassword: '', + UptimeKumaMonitorScope: optionMap.UptimeKumaMonitorScope ?? 'all', + UptimeKumaSelectedSites: optionMap.UptimeKumaSelectedSites ?? '', + UptimeKumaSyncInterval: optionMap.UptimeKumaSyncInterval ?? '5', + UptimeKumaInterval: optionMap.UptimeKumaInterval ?? '60', + UptimeKumaRetry: optionMap.UptimeKumaRetry ?? '0', + UptimeKumaRetryInterval: optionMap.UptimeKumaRetryInterval ?? '60', + UptimeKumaTimeout: optionMap.UptimeKumaTimeout ?? '48', ServerAddress: resolvedServerAddress, }); @@ -656,6 +681,64 @@ export function SettingsPage() { }); }; + const handleUptimeKumaSave = () => { + void runBusyAction('uptimekuma-save', async () => { + const syncInt = Number.parseInt(operationFields.UptimeKumaSyncInterval, 10); + const interval = Number.parseInt(operationFields.UptimeKumaInterval, 10); + const retry = Number.parseInt(operationFields.UptimeKumaRetry, 10); + const retryInt = Number.parseInt(operationFields.UptimeKumaRetryInterval, 10); + const timeout = Number.parseInt(operationFields.UptimeKumaTimeout, 10); + + if (operationFields.UptimeKumaEnabled) { + if (!operationFields.UptimeKumaUrl.trim()) { + throw new Error('请输入 Uptime Kuma 地址。'); + } + if (!operationFields.UptimeKumaUsername.trim()) { + throw new Error('请输入 Uptime Kuma 用户名。'); + } + } + if (Number.isNaN(syncInt) || syncInt <= 0) { + throw new Error('同步间隔必须为正整数。'); + } + if (Number.isNaN(interval) || interval <= 0) { + throw new Error('心跳间隔必须为正整数。'); + } + if (Number.isNaN(retry) || retry < 0) { + throw new Error('重试次数必须为非负整数。'); + } + if (Number.isNaN(retryInt) || retryInt <= 0) { + throw new Error('心跳重试间隔必须为正整数。'); + } + if (Number.isNaN(timeout) || timeout <= 0) { + throw new Error('请求超时必须为正整数。'); + } + + await saveOptionEntries( + [ + ['UptimeKumaEnabled', String(operationFields.UptimeKumaEnabled)], + ['UptimeKumaUrl', operationFields.UptimeKumaUrl.trim()], + ['UptimeKumaUsername', operationFields.UptimeKumaUsername.trim()], + ['UptimeKumaPassword', operationFields.UptimeKumaPassword], + ['UptimeKumaMonitorScope', operationFields.UptimeKumaMonitorScope], + ['UptimeKumaSelectedSites', operationFields.UptimeKumaSelectedSites], + ['UptimeKumaSyncInterval', String(syncInt)], + ['UptimeKumaInterval', String(interval)], + ['UptimeKumaRetry', String(retry)], + ['UptimeKumaRetryInterval', String(retryInt)], + ['UptimeKumaTimeout', String(timeout)], + ], + 'Uptime Kuma 设置已保存。', + ); + }); + }; + + const handleUptimeKumaSync = () => { + void runBusyAction('uptimekuma-sync', async () => { + await syncUptimeKuma(); + setFeedback({ tone: 'success', message: '同步任务已成功执行!' }); + }); + }; + const renderTabContent = () => { if (profileQuery.isLoading || publicStatusQuery.isLoading) { return ; @@ -1307,6 +1390,198 @@ export function SettingsPage() { )} + + + + {busyKey === 'uptimekuma-sync' ? '同步中...' : '立即同步'} + + + {busyKey === 'uptimekuma-save' ? '保存中...' : '保存设置'} + + + } + > +
+ + setOperationFields((previous) => ({ + ...previous, + UptimeKumaEnabled: checked, + })) + } + /> + + {operationFields.UptimeKumaEnabled ? ( +
+
+ + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaUrl: event.target.value, + })) + } + placeholder="http://localhost:3001" + /> + + + + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaUsername: event.target.value, + })) + } + placeholder="请输入用户名" + /> + +
+ +
+ + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaPassword: event.target.value, + })) + } + placeholder="请输入密码(留空表示不更新)" + /> + + + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaSyncInterval: event.target.value, + })) + } + /> + + + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaMonitorScope: event.target.value, + })) + } + > + + + + + + +
+ + {operationFields.UptimeKumaMonitorScope === 'selected' ? ( +
+
+ 已选站点 + setUptimeKumaModalOpen(true)}> + 选择监控站点 + +
+
+ {operationFields.UptimeKumaSelectedSites + ? operationFields.UptimeKumaSelectedSites.split(',').join(', ') + : '未选择任何站点,同步不会执行。'} +
+
+ ) : null} + +
+

+ Uptime Kuma 属性配置 +

+
+ + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaInterval: event.target.value, + })) + } + /> + + + + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaRetry: event.target.value, + })) + } + /> + + + + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaRetryInterval: event.target.value, + })) + } + /> + + + + + setOperationFields((previous) => ({ + ...previous, + UptimeKumaTimeout: event.target.value, + })) + } + /> + +
+
+
+ ) : null} +
+
+
@@ -2050,6 +2325,22 @@ export function SettingsPage() { }} /> + setUptimeKumaModalOpen(false)} + onSave={(sites) => + setOperationFields((previous) => ({ + ...previous, + UptimeKumaSelectedSites: sites.join(','), + })) + } + /> + void; + onSave: (sites: string[]) => void; +}) { + const [searchTerm, setSearchTerm] = useState(''); + const [tempSelected, setTempSelected] = useState>(new Set()); + + const { data: routes = [], isLoading, error } = useQuery({ + queryKey: ['proxy-routes'], + queryFn: getProxyRoutes, + enabled: isOpen, + }); + + useEffect(() => { + if (isOpen) { + setTempSelected(new Set(selectedSites.map(s => s.trim()).filter(Boolean))); + setSearchTerm(''); + } + }, [isOpen, selectedSites]); + + const toggleSite = (siteName: string) => { + setTempSelected((prev) => { + const next = new Set(prev); + if (next.has(siteName)) { + next.delete(siteName); + } else { + next.add(siteName); + } + return next; + }); + }; + + const handleSelectAll = () => { + setTempSelected((prev) => { + const next = new Set(prev); + filteredRoutes.forEach((route) => { + next.add(route.site_name); + }); + return next; + }); + }; + + const handleDeselectAll = () => { + setTempSelected((prev) => { + const next = new Set(prev); + filteredRoutes.forEach((route) => { + next.delete(route.site_name); + }); + return next; + }); + }; + + const handleSave = () => { + onSave(Array.from(tempSelected)); + onClose(); + }; + + const filteredRoutes = routes.filter( + (route) => + route.site_name.toLowerCase().includes(searchTerm.toLowerCase()) || + route.primary_domain.toLowerCase().includes(searchTerm.toLowerCase()), + ); + + return ( + + + 取消 + + + 保存选择 + +
+ } + > +
+
+
+ + setSearchTerm(event.target.value)} + placeholder="按名称或域名搜索..." + /> + +
+
+ + 全选过滤项 + + + 清空过滤项 + +
+
+ + {isLoading ? : null} + + {error ? ( + + ) : null} + + {!isLoading && !error && routes.length === 0 ? ( +
+ 暂无可用的代理站点。 +
+ ) : null} + + {!isLoading && !error && routes.length > 0 ? ( +
+ + + + + + + + + + + {filteredRoutes.map((route) => { + const isChecked = tempSelected.has(route.site_name); + return ( + toggleSite(route.site_name)} + className="cursor-pointer hover:bg-[var(--surface-elevated)]" + > + + + + + + ); + })} + {filteredRoutes.length === 0 ? ( + + + + ) : null} + +
选择站点名称主域名状态
e.stopPropagation()}> + toggleSite(route.site_name)} + className="h-4 w-4 rounded border-gray-300 text-indigo-600 focus:ring-indigo-500" + /> + + {route.site_name} + + {route.primary_domain} + + + {route.enabled ? '启用' : '禁用'} + +
+ 无匹配的站点 +
+
+ ) : null} + +
+ 已选择 {tempSelected.size} 个监控站点 +
+
+ + ); +}