From 2ac083a22500197436e4e94e7a18086bbd335ab9 Mon Sep 17 00:00:00 2001 From: ryan Date: Mon, 9 Mar 2026 23:05:35 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E4=BB=A3=E7=90=86=E8=8A=82?= =?UTF-8?q?=E7=82=B9=E5=8A=9F=E8=83=BD=EF=BC=8C=E5=8C=85=E6=8B=AC=E8=8A=82?= =?UTF-8?q?=E7=82=B9=E6=B3=A8=E5=86=8C=E3=80=81=E5=BF=83=E8=B7=B3=E6=A3=80?= =?UTF-8?q?=E6=B5=8B=E5=92=8C=E9=85=8D=E7=BD=AE=E8=8E=B7=E5=8F=96=EF=BC=8C?= =?UTF-8?q?=E5=AE=8C=E5=96=84=E7=9B=B8=E5=85=B3API=E8=B7=AF=E7=94=B1?= =?UTF-8?q?=E5=92=8C=E6=9C=8D=E5=8A=A1=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- atsf_agent/cmd/agent/main.go | 3 + atsf_agent/go.mod | 3 + atsf_agent/internal/protocol/agent_api.go | 25 +++ atsf_server/common/constants.go | 2 + atsf_server/common/init.go | 3 + atsf_server/controller/agent.go | 128 ++++++++++++ atsf_server/middleware/agent-auth.go | 30 +++ atsf_server/model/apply_log.go | 27 +++ atsf_server/model/main.go | 8 + atsf_server/model/node.go | 29 +++ atsf_server/router/api-router.go | 18 ++ atsf_server/router/api_phase2_test.go | 176 +++++++++++++++++ atsf_server/service/agent.go | 225 ++++++++++++++++++++++ 13 files changed, 677 insertions(+) create mode 100644 atsf_agent/cmd/agent/main.go create mode 100644 atsf_agent/go.mod create mode 100644 atsf_agent/internal/protocol/agent_api.go create mode 100644 atsf_server/controller/agent.go create mode 100644 atsf_server/middleware/agent-auth.go create mode 100644 atsf_server/model/apply_log.go create mode 100644 atsf_server/model/node.go create mode 100644 atsf_server/router/api_phase2_test.go create mode 100644 atsf_server/service/agent.go diff --git a/atsf_agent/cmd/agent/main.go b/atsf_agent/cmd/agent/main.go new file mode 100644 index 00000000..38dd16da --- /dev/null +++ b/atsf_agent/cmd/agent/main.go @@ -0,0 +1,3 @@ +package main + +func main() {} diff --git a/atsf_agent/go.mod b/atsf_agent/go.mod new file mode 100644 index 00000000..4462ea22 --- /dev/null +++ b/atsf_agent/go.mod @@ -0,0 +1,3 @@ +module atsflare-agent + +go 1.18 diff --git a/atsf_agent/internal/protocol/agent_api.go b/atsf_agent/internal/protocol/agent_api.go new file mode 100644 index 00000000..5cc8ffa3 --- /dev/null +++ b/atsf_agent/internal/protocol/agent_api.go @@ -0,0 +1,25 @@ +package protocol + +type NodePayload struct { + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + AgentVersion string `json:"agent_version"` + NginxVersion string `json:"nginx_version"` + CurrentVersion string `json:"current_version"` + LastError string `json:"last_error"` +} + +type ApplyLogPayload struct { + NodeID string `json:"node_id"` + Version string `json:"version"` + Result string `json:"result"` + Message string `json:"message"` +} + +type ActiveConfigResponse struct { + Version string `json:"version"` + Checksum string `json:"checksum"` + RenderedConfig string `json:"rendered_config"` + CreatedAt string `json:"created_at"` +} diff --git a/atsf_server/common/constants.go b/atsf_server/common/constants.go index 14376820..51fc579f 100644 --- a/atsf_server/common/constants.go +++ b/atsf_server/common/constants.go @@ -45,6 +45,8 @@ var WeChatAccountQRCodeImageURL = "" var TurnstileSiteKey = "" var TurnstileSecretKey = "" +var AgentToken = "" +var NodeOfflineThreshold = 2 * time.Minute const ( RoleGuestUser = 0 diff --git a/atsf_server/common/init.go b/atsf_server/common/init.go index 50e7f05e..c66090ca 100644 --- a/atsf_server/common/init.go +++ b/atsf_server/common/init.go @@ -50,6 +50,9 @@ func init() { if os.Getenv("UPLOAD_PATH") != "" { UploadPath = os.Getenv("UPLOAD_PATH") } + if os.Getenv("AGENT_TOKEN") != "" { + AgentToken = os.Getenv("AGENT_TOKEN") + } if *LogDir != "" { var err error *LogDir, err = filepath.Abs(*LogDir) diff --git a/atsf_server/controller/agent.go b/atsf_server/controller/agent.go new file mode 100644 index 00000000..4863711f --- /dev/null +++ b/atsf_server/controller/agent.go @@ -0,0 +1,128 @@ +package controller + +import ( + "encoding/json" + "gin-template/service" + "github.com/gin-gonic/gin" + "net/http" +) + +func AgentRegister(c *gin.Context) { + var payload service.AgentNodePayload + if err := json.NewDecoder(c.Request.Body).Decode(&payload); err != nil { + c.JSON(http.StatusBadRequest, gin.H{ + "success": false, + "message": "无效的参数", + }) + return + } + node, err := service.RegisterNode(payload) + if err != nil { + c.JSON(http.StatusOK, gin.H{ + "success": false, + "message": err.Error(), + }) + return + } + c.JSON(http.StatusOK, gin.H{ + "success": true, + "message": "", + "data": node, + }) +} + +func AgentHeartbeat(c *gin.Context) { + var payload service.AgentNodePayload + if err := json.NewDecoder(c.Request.Body).Decode(&payload); err != nil { + c.JSON(http.StatusBadRequest, gin.H{ + "success": false, + "message": "无效的参数", + }) + return + } + node, err := service.HeartbeatNode(payload) + if err != nil { + c.JSON(http.StatusOK, gin.H{ + "success": false, + "message": err.Error(), + }) + return + } + c.JSON(http.StatusOK, gin.H{ + "success": true, + "message": "", + "data": node, + }) +} + +func AgentGetActiveConfig(c *gin.Context) { + config, err := service.GetActiveConfigForAgent() + if err != nil { + c.JSON(http.StatusOK, gin.H{ + "success": false, + "message": "当前没有激活版本", + }) + return + } + c.JSON(http.StatusOK, gin.H{ + "success": true, + "message": "", + "data": config, + }) +} + +func AgentReportApplyLog(c *gin.Context) { + var payload service.ApplyLogPayload + if err := json.NewDecoder(c.Request.Body).Decode(&payload); err != nil { + c.JSON(http.StatusBadRequest, gin.H{ + "success": false, + "message": "无效的参数", + }) + return + } + log, err := service.ReportApplyLog(payload) + if err != nil { + c.JSON(http.StatusOK, gin.H{ + "success": false, + "message": err.Error(), + }) + return + } + c.JSON(http.StatusOK, gin.H{ + "success": true, + "message": "", + "data": log, + }) +} + +func GetNodes(c *gin.Context) { + nodes, err := service.ListNodeViews() + if err != nil { + c.JSON(http.StatusOK, gin.H{ + "success": false, + "message": err.Error(), + }) + return + } + c.JSON(http.StatusOK, gin.H{ + "success": true, + "message": "", + "data": nodes, + }) +} + +func GetApplyLogs(c *gin.Context) { + logs, err := service.ListApplyLogs(c.Query("node_id")) + if err != nil { + c.JSON(http.StatusOK, gin.H{ + "success": false, + "message": err.Error(), + }) + return + } + c.JSON(http.StatusOK, gin.H{ + "success": true, + "message": "", + "data": logs, + }) +} diff --git a/atsf_server/middleware/agent-auth.go b/atsf_server/middleware/agent-auth.go new file mode 100644 index 00000000..4cbfd15c --- /dev/null +++ b/atsf_server/middleware/agent-auth.go @@ -0,0 +1,30 @@ +package middleware + +import ( + "gin-template/common" + "github.com/gin-gonic/gin" + "net/http" +) + +func AgentAuth() func(c *gin.Context) { + return func(c *gin.Context) { + token := c.GetHeader("X-Agent-Token") + if common.AgentToken == "" { + c.JSON(http.StatusUnauthorized, gin.H{ + "success": false, + "message": "Agent Token 未配置", + }) + c.Abort() + return + } + if token == "" || token != common.AgentToken { + c.JSON(http.StatusUnauthorized, gin.H{ + "success": false, + "message": "无权进行此操作,Agent Token 无效", + }) + c.Abort() + return + } + c.Next() + } +} diff --git a/atsf_server/model/apply_log.go b/atsf_server/model/apply_log.go new file mode 100644 index 00000000..888af1e0 --- /dev/null +++ b/atsf_server/model/apply_log.go @@ -0,0 +1,27 @@ +package model + +import "time" + +type ApplyLog struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"index;size:64;not null"` + Version string `json:"version" gorm:"size:32;not null"` + Result string `json:"result" gorm:"size:32;not null"` + Message string `json:"message" gorm:"size:1024"` + CreatedAt time.Time `json:"created_at"` +} + +func ListApplyLogs(nodeID string) (logs []*ApplyLog, err error) { + query := DB.Order("id desc") + if nodeID != "" { + query = query.Where("node_id = ?", nodeID) + } + err = query.Find(&logs).Error + return logs, err +} + +func GetLatestApplyLog(nodeID string) (*ApplyLog, error) { + log := &ApplyLog{} + err := DB.Where("node_id = ?", nodeID).Order("id desc").First(log).Error + return log, err +} diff --git a/atsf_server/model/main.go b/atsf_server/model/main.go index 25e80ff6..412d94b2 100644 --- a/atsf_server/model/main.go +++ b/atsf_server/model/main.go @@ -72,6 +72,14 @@ func InitDB() (err error) { if err != nil { return err } + err = db.AutoMigrate(&Node{}) + if err != nil { + return err + } + err = db.AutoMigrate(&ApplyLog{}) + if err != nil { + return err + } err = createRootAccountIfNeed() return err } else { diff --git a/atsf_server/model/node.go b/atsf_server/model/node.go new file mode 100644 index 00000000..1c1399cf --- /dev/null +++ b/atsf_server/model/node.go @@ -0,0 +1,29 @@ +package model + +import "time" + +type Node struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"uniqueIndex;size:64;not null"` + Name string `json:"name" gorm:"size:128;not null"` + IP string `json:"ip" gorm:"size:64;not null"` + AgentVersion string `json:"agent_version" gorm:"size:64;not null"` + NginxVersion string `json:"nginx_version" gorm:"size:64"` + Status string `json:"status" gorm:"size:16;not null;default:'offline'"` + CurrentVersion string `json:"current_version" gorm:"size:32"` + LastSeenAt time.Time `json:"last_seen_at"` + LastError string `json:"last_error" gorm:"size:1024"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +func ListNodes() (nodes []*Node, err error) { + err = DB.Order("id desc").Find(&nodes).Error + return nodes, err +} + +func GetNodeByNodeID(nodeID string) (*Node, error) { + node := &Node{} + err := DB.Where("node_id = ?", nodeID).First(node).Error + return node, err +} diff --git a/atsf_server/router/api-router.go b/atsf_server/router/api-router.go index f1e93836..095d2145 100644 --- a/atsf_server/router/api-router.go +++ b/atsf_server/router/api-router.go @@ -78,5 +78,23 @@ func SetApiRouter(router *gin.Engine) { configVersionRoute.POST("/publish", controller.PublishConfigVersion) configVersionRoute.PUT("/:id/activate", controller.ActivateConfigVersion) } + nodeRoute := apiRouter.Group("/nodes") + nodeRoute.Use(middleware.AdminAuth()) + { + nodeRoute.GET("/", controller.GetNodes) + } + applyLogRoute := apiRouter.Group("/apply-logs") + applyLogRoute.Use(middleware.AdminAuth()) + { + applyLogRoute.GET("/", controller.GetApplyLogs) + } + agentRoute := apiRouter.Group("/agent") + agentRoute.Use(middleware.AgentAuth()) + { + agentRoute.POST("/nodes/register", controller.AgentRegister) + agentRoute.POST("/nodes/heartbeat", controller.AgentHeartbeat) + agentRoute.GET("/config-versions/active", controller.AgentGetActiveConfig) + agentRoute.POST("/apply-logs", controller.AgentReportApplyLog) + } } } diff --git a/atsf_server/router/api_phase2_test.go b/atsf_server/router/api_phase2_test.go new file mode 100644 index 00000000..9d17f95c --- /dev/null +++ b/atsf_server/router/api_phase2_test.go @@ -0,0 +1,176 @@ +package router_test + +import ( + "bytes" + "encoding/json" + "gin-template/common" + "gin-template/model" + "gin-template/router" + "gin-template/service" + "github.com/gin-contrib/sessions" + "github.com/gin-contrib/sessions/cookie" + "github.com/gin-gonic/gin" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestPhase2AgentLifecycle(t *testing.T) { + gin.SetMode(gin.TestMode) + common.RedisEnabled = false + common.AgentToken = "phase2-agent-token" + setupTestDB(t) + + engine := gin.New() + engine.Use(sessions.Sessions("session", cookie.NewStore([]byte("test-secret")))) + router.SetApiRouter(engine) + + adminToken := prepareRootToken(t) + + createRouteAndPublishVersion(t, engine, adminToken) + + unauthorizedRequest := httptest.NewRequest(http.MethodPost, "/api/agent/nodes/register", bytes.NewReader([]byte(`{}`))) + unauthorizedRecorder := httptest.NewRecorder() + engine.ServeHTTP(unauthorizedRecorder, unauthorizedRequest) + if unauthorizedRecorder.Code != http.StatusUnauthorized { + t.Fatalf("expected unauthorized status for missing agent token, got %d", unauthorizedRecorder.Code) + } + + nodePayload := map[string]any{ + "node_id": "node-001", + "name": "shanghai-edge-1", + "ip": "10.0.0.8", + "agent_version": "0.1.0", + "nginx_version": "1.25.5", + "current_version": "", + "last_error": "", + } + resp := performAgentJSONRequest(t, engine, http.MethodPost, "/api/agent/nodes/register", nodePayload) + var registeredNode model.Node + decodeResponseData(t, resp, ®isteredNode) + if registeredNode.NodeID != "node-001" || registeredNode.Status != service.NodeStatusOnline { + t.Fatal("expected node registration to persist online node state") + } + + heartbeatPayload := map[string]any{ + "node_id": "node-001", + "name": "shanghai-edge-1", + "ip": "10.0.0.9", + "agent_version": "0.1.1", + "nginx_version": "1.25.5", + "current_version": "", + "last_error": "", + } + resp = performAgentJSONRequest(t, engine, http.MethodPost, "/api/agent/nodes/heartbeat", heartbeatPayload) + decodeResponseData(t, resp, ®isteredNode) + if registeredNode.IP != "10.0.0.9" || registeredNode.AgentVersion != "0.1.1" { + t.Fatal("expected heartbeat to update node metadata") + } + + activeConfigResp := performAgentJSONRequest(t, engine, http.MethodGet, "/api/agent/config-versions/active", nil) + var activeConfig service.AgentConfigResponse + decodeResponseData(t, activeConfigResp, &activeConfig) + if activeConfig.Version == "" || activeConfig.RenderedConfig == "" || activeConfig.Checksum == "" { + t.Fatal("expected active config response to contain version payload") + } + + successApplyResp := performAgentJSONRequest(t, engine, http.MethodPost, "/api/agent/apply-logs", map[string]any{ + "node_id": "node-001", + "version": activeConfig.Version, + "result": service.ApplyResultOK, + "message": "apply ok", + }) + var successApplyLog model.ApplyLog + decodeResponseData(t, successApplyResp, &successApplyLog) + if successApplyLog.Result != service.ApplyResultOK { + t.Fatal("expected apply log success to be recorded") + } + + failedApplyResp := performAgentJSONRequest(t, engine, http.MethodPost, "/api/agent/apply-logs", map[string]any{ + "node_id": "node-001", + "version": activeConfig.Version, + "result": service.ApplyResultFailed, + "message": "nginx reload failed", + }) + var failedApplyLog model.ApplyLog + decodeResponseData(t, failedApplyResp, &failedApplyLog) + if failedApplyLog.Result != service.ApplyResultFailed { + t.Fatal("expected failed apply log to be recorded") + } + + nodesResp := performJSONRequest(t, engine, adminToken, http.MethodGet, "/api/nodes/", nil) + var nodes []service.NodeView + decodeResponseData(t, nodesResp, &nodes) + if len(nodes) != 1 { + t.Fatalf("expected 1 node, got %d", len(nodes)) + } + if nodes[0].LatestApplyResult != service.ApplyResultFailed || nodes[0].LatestApplyMessage != "nginx reload failed" { + t.Fatal("expected node list to expose latest apply status") + } + if nodes[0].CurrentVersion != activeConfig.Version { + t.Fatal("expected node current_version to remain at last successful version") + } + if nodes[0].LastError != "nginx reload failed" { + t.Fatal("expected node last_error to reflect failed apply") + } + + logsResp := performJSONRequest(t, engine, adminToken, http.MethodGet, "/api/apply-logs/?node_id=node-001", nil) + var logs []model.ApplyLog + decodeResponseData(t, logsResp, &logs) + if len(logs) != 2 { + t.Fatalf("expected 2 apply logs, got %d", len(logs)) + } + + oldTime := time.Now().Add(-common.NodeOfflineThreshold - time.Minute) + if err := model.DB.Model(&model.Node{}).Where("node_id = ?", "node-001").Update("last_seen_at", oldTime).Error; err != nil { + t.Fatalf("failed to update node last_seen_at: %v", err) + } + nodesResp = performJSONRequest(t, engine, adminToken, http.MethodGet, "/api/nodes/", nil) + decodeResponseData(t, nodesResp, &nodes) + if nodes[0].Status != service.NodeStatusOffline { + t.Fatal("expected node to be shown as offline after timeout") + } +} + +func performAgentJSONRequest(t *testing.T, engine http.Handler, method string, path string, body any) apiResponse { + t.Helper() + var payload []byte + var err error + if body != nil { + payload, err = json.Marshal(body) + if err != nil { + t.Fatalf("failed to marshal request body: %v", err) + } + } + req := httptest.NewRequest(method, path, bytes.NewReader(payload)) + if body != nil { + req.Header.Set("Content-Type", "application/json") + } + req.Header.Set("X-Agent-Token", common.AgentToken) + recorder := httptest.NewRecorder() + engine.ServeHTTP(recorder, req) + if recorder.Code != http.StatusOK { + t.Fatalf("unexpected status %d for %s %s: %s", recorder.Code, method, path, recorder.Body.String()) + } + var resp apiResponse + if err = json.Unmarshal(recorder.Body.Bytes(), &resp); err != nil { + t.Fatalf("failed to unmarshal response: %v", err) + } + if !resp.Success { + t.Fatalf("request %s %s failed: %s", method, path, resp.Message) + } + return resp +} + +func createRouteAndPublishVersion(t *testing.T, engine http.Handler, adminToken string) { + t.Helper() + createBody := map[string]any{ + "domain": "agent.example.com", + "origin_url": "https://agent-origin.internal", + "enabled": true, + "remark": "agent route", + } + performJSONRequest(t, engine, adminToken, http.MethodPost, "/api/proxy-routes/", createBody) + performJSONRequest(t, engine, adminToken, http.MethodPost, "/api/config-versions/publish", nil) +} diff --git a/atsf_server/service/agent.go b/atsf_server/service/agent.go new file mode 100644 index 00000000..ae073ef9 --- /dev/null +++ b/atsf_server/service/agent.go @@ -0,0 +1,225 @@ +package service + +import ( + "errors" + "gin-template/common" + "gin-template/model" + "strings" + "time" + + "gorm.io/gorm" +) + +const ( + NodeStatusOnline = "online" + NodeStatusOffline = "offline" + ApplyResultOK = "success" + ApplyResultFailed = "failed" +) + +type AgentNodePayload struct { + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + AgentVersion string `json:"agent_version"` + NginxVersion string `json:"nginx_version"` + CurrentVersion string `json:"current_version"` + LastError string `json:"last_error"` +} + +type ApplyLogPayload struct { + NodeID string `json:"node_id"` + Version string `json:"version"` + Result string `json:"result"` + Message string `json:"message"` +} + +type AgentConfigResponse struct { + Version string `json:"version"` + Checksum string `json:"checksum"` + RenderedConfig string `json:"rendered_config"` + CreatedAt time.Time `json:"created_at"` +} + +type NodeView struct { + ID uint `json:"id"` + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + AgentVersion string `json:"agent_version"` + NginxVersion string `json:"nginx_version"` + Status string `json:"status"` + CurrentVersion string `json:"current_version"` + LastSeenAt time.Time `json:"last_seen_at"` + LastError string `json:"last_error"` + LatestApplyResult string `json:"latest_apply_result"` + LatestApplyMessage string `json:"latest_apply_message"` + LatestApplyAt *time.Time `json:"latest_apply_at"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +func RegisterNode(payload AgentNodePayload) (*model.Node, error) { + return upsertNode(payload) +} + +func HeartbeatNode(payload AgentNodePayload) (*model.Node, error) { + return upsertNode(payload) +} + +func GetActiveConfigForAgent() (*AgentConfigResponse, error) { + version, err := model.GetActiveConfigVersion() + if err != nil { + return nil, err + } + return &AgentConfigResponse{ + Version: version.Version, + Checksum: version.Checksum, + RenderedConfig: version.RenderedConfig, + CreatedAt: version.CreatedAt, + }, nil +} + +func ReportApplyLog(payload ApplyLogPayload) (*model.ApplyLog, error) { + now := time.Now() + payload.NodeID = strings.TrimSpace(payload.NodeID) + payload.Version = strings.TrimSpace(payload.Version) + payload.Result = strings.TrimSpace(strings.ToLower(payload.Result)) + payload.Message = strings.TrimSpace(payload.Message) + if payload.NodeID == "" { + return nil, errors.New("node_id 不能为空") + } + if payload.Version == "" { + return nil, errors.New("version 不能为空") + } + if payload.Result != ApplyResultOK && payload.Result != ApplyResultFailed { + return nil, errors.New("result 仅支持 success 或 failed") + } + + log := &model.ApplyLog{ + NodeID: payload.NodeID, + Version: payload.Version, + Result: payload.Result, + Message: payload.Message, + CreatedAt: now, + } + err := model.DB.Transaction(func(tx *gorm.DB) error { + node := &model.Node{} + if err := tx.Where("node_id = ?", payload.NodeID).First(node).Error; err != nil { + return err + } + node.Status = NodeStatusOnline + node.LastSeenAt = now + if payload.Result == ApplyResultOK { + node.CurrentVersion = payload.Version + node.LastError = "" + } else { + node.LastError = payload.Message + } + if err := tx.Create(log).Error; err != nil { + return err + } + return tx.Model(node).Select("status", "last_seen_at", "current_version", "last_error").Updates(node).Error + }) + if err != nil { + return nil, err + } + return log, nil +} + +func ListNodeViews() ([]*NodeView, error) { + nodes, err := model.ListNodes() + if err != nil { + return nil, err + } + views := make([]*NodeView, 0, len(nodes)) + for _, node := range nodes { + view := &NodeView{ + ID: node.ID, + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + AgentVersion: node.AgentVersion, + NginxVersion: node.NginxVersion, + Status: computeNodeStatus(node.LastSeenAt), + CurrentVersion: node.CurrentVersion, + LastSeenAt: node.LastSeenAt, + LastError: node.LastError, + CreatedAt: node.CreatedAt, + UpdatedAt: node.UpdatedAt, + } + if log, err := model.GetLatestApplyLog(node.NodeID); err == nil { + view.LatestApplyResult = log.Result + view.LatestApplyMessage = log.Message + view.LatestApplyAt = &log.CreatedAt + } + views = append(views, view) + } + return views, nil +} + +func ListApplyLogs(nodeID string) ([]*model.ApplyLog, error) { + return model.ListApplyLogs(strings.TrimSpace(nodeID)) +} + +func upsertNode(payload AgentNodePayload) (*model.Node, error) { + now := time.Now() + payload.NodeID = strings.TrimSpace(payload.NodeID) + payload.Name = strings.TrimSpace(payload.Name) + payload.IP = strings.TrimSpace(payload.IP) + payload.AgentVersion = strings.TrimSpace(payload.AgentVersion) + payload.NginxVersion = strings.TrimSpace(payload.NginxVersion) + payload.CurrentVersion = strings.TrimSpace(payload.CurrentVersion) + payload.LastError = strings.TrimSpace(payload.LastError) + if payload.NodeID == "" { + return nil, errors.New("node_id 不能为空") + } + if payload.Name == "" { + return nil, errors.New("name 不能为空") + } + if payload.IP == "" { + return nil, errors.New("ip 不能为空") + } + if payload.AgentVersion == "" { + return nil, errors.New("agent_version 不能为空") + } + + node := &model.Node{} + err := model.DB.Where("node_id = ?", payload.NodeID).First(node).Error + if err != nil { + if !errors.Is(err, gorm.ErrRecordNotFound) { + return nil, err + } + node = &model.Node{ + NodeID: payload.NodeID, + } + } + node.Name = payload.Name + node.IP = payload.IP + node.AgentVersion = payload.AgentVersion + node.NginxVersion = payload.NginxVersion + node.Status = NodeStatusOnline + node.CurrentVersion = payload.CurrentVersion + node.LastSeenAt = now + node.LastError = payload.LastError + if node.ID == 0 { + if err = model.DB.Create(node).Error; err != nil { + return nil, err + } + return node, nil + } + if err = model.DB.Model(node).Select("name", "ip", "agent_version", "nginx_version", "status", "current_version", "last_seen_at", "last_error").Updates(node).Error; err != nil { + return nil, err + } + return node, nil +} + +func computeNodeStatus(lastSeenAt time.Time) string { + if lastSeenAt.IsZero() { + return NodeStatusOffline + } + if time.Since(lastSeenAt) > common.NodeOfflineThreshold { + return NodeStatusOffline + } + return NodeStatusOnline +}