mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-10 17:26:38 +08:00
添加代理节点功能,包括节点注册、心跳检测和配置获取,完善相关API路由和服务逻辑
This commit is contained in:
@@ -0,0 +1,3 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
func main() {}
|
||||||
@@ -0,0 +1,3 @@
|
|||||||
|
module atsflare-agent
|
||||||
|
|
||||||
|
go 1.18
|
||||||
@@ -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"`
|
||||||
|
}
|
||||||
@@ -45,6 +45,8 @@ var WeChatAccountQRCodeImageURL = ""
|
|||||||
|
|
||||||
var TurnstileSiteKey = ""
|
var TurnstileSiteKey = ""
|
||||||
var TurnstileSecretKey = ""
|
var TurnstileSecretKey = ""
|
||||||
|
var AgentToken = ""
|
||||||
|
var NodeOfflineThreshold = 2 * time.Minute
|
||||||
|
|
||||||
const (
|
const (
|
||||||
RoleGuestUser = 0
|
RoleGuestUser = 0
|
||||||
|
|||||||
@@ -50,6 +50,9 @@ func init() {
|
|||||||
if os.Getenv("UPLOAD_PATH") != "" {
|
if os.Getenv("UPLOAD_PATH") != "" {
|
||||||
UploadPath = os.Getenv("UPLOAD_PATH")
|
UploadPath = os.Getenv("UPLOAD_PATH")
|
||||||
}
|
}
|
||||||
|
if os.Getenv("AGENT_TOKEN") != "" {
|
||||||
|
AgentToken = os.Getenv("AGENT_TOKEN")
|
||||||
|
}
|
||||||
if *LogDir != "" {
|
if *LogDir != "" {
|
||||||
var err error
|
var err error
|
||||||
*LogDir, err = filepath.Abs(*LogDir)
|
*LogDir, err = filepath.Abs(*LogDir)
|
||||||
|
|||||||
@@ -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,
|
||||||
|
})
|
||||||
|
}
|
||||||
@@ -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()
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -72,6 +72,14 @@ func InitDB() (err error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
err = db.AutoMigrate(&Node{})
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
err = db.AutoMigrate(&ApplyLog{})
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
err = createRootAccountIfNeed()
|
err = createRootAccountIfNeed()
|
||||||
return err
|
return err
|
||||||
} else {
|
} else {
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -78,5 +78,23 @@ func SetApiRouter(router *gin.Engine) {
|
|||||||
configVersionRoute.POST("/publish", controller.PublishConfigVersion)
|
configVersionRoute.POST("/publish", controller.PublishConfigVersion)
|
||||||
configVersionRoute.PUT("/:id/activate", controller.ActivateConfigVersion)
|
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)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user