From 8c481985c2b96baf7a5d193812e0b1fab4d21433 Mon Sep 17 00:00:00 2001 From: Buri Date: Fri, 4 Sep 2026 18:17:52 +0800 Subject: [PATCH] feat: enforce telemetry-only mode and safe agent decommission lifecycle --- cmd/dashboard/controller/controller.go | 1 + cmd/dashboard/controller/fm.go | 4 +++ cmd/dashboard/controller/server.go | 48 ++++++++++++++++++++++++++ cmd/dashboard/controller/terminal.go | 4 +++ model/server.go | 38 ++++++++++++++++++++ service/rpc/auth.go | 2 +- service/rpc/io_stream_rpc.go | 4 +++ service/rpc/nezha.go | 21 +++++++++++ 8 files changed, 121 insertions(+), 1 deletion(-) diff --git a/cmd/dashboard/controller/controller.go b/cmd/dashboard/controller/controller.go index 7927ca4d..9ea99633 100644 --- a/cmd/dashboard/controller/controller.go +++ b/cmd/dashboard/controller/controller.go @@ -132,6 +132,7 @@ func routers(r *gin.Engine, frontendDist fs.FS) { auth.POST("/batch-delete/server", restScopeMiddleware(model.ScopeInventoryDelete), commonHandler(batchDeleteServer)) auth.POST("/batch-move/server", restScopeMiddleware(model.ScopeServerWrite), commonHandler(batchMoveServer)) auth.POST("/force-update/server", restScopeMiddleware(model.ScopeServerWrite), commonHandler(forceUpdateServer)) + auth.POST("/batch-lockdown/server", restScopeMiddleware(model.ScopeServerWrite), commonHandler(batchLockdownOnlineServers)) auth.POST("/server-group", restScopeMiddleware(model.ScopeServerWrite), commonHandler(createServerGroup)) auth.PATCH("/server-group/:id", restScopeMiddleware(model.ScopeServerWrite), commonHandler(updateServerGroup)) auth.POST("/batch-delete/server-group", restScopeMiddleware(model.ScopeInventoryDelete), commonHandler(batchDeleteServerGroup)) diff --git a/cmd/dashboard/controller/fm.go b/cmd/dashboard/controller/fm.go index dfdb2728..bd9216fa 100644 --- a/cmd/dashboard/controller/fm.go +++ b/cmd/dashboard/controller/fm.go @@ -44,6 +44,10 @@ func createFM(c *gin.Context) (*model.CreateFMResponse, error) { return nil, singleton.Localizer.ErrorT("permission denied") } + if server.IsTelemetryOnly() { + return nil, singleton.Localizer.ErrorT("file manager is disabled: server is in telemetry-only mode") + } + streamId, err := uuid.GenerateUUID() if err != nil { return nil, err diff --git a/cmd/dashboard/controller/server.go b/cmd/dashboard/controller/server.go index c5c6a4ac..62e8ee95 100644 --- a/cmd/dashboard/controller/server.go +++ b/cmd/dashboard/controller/server.go @@ -556,3 +556,51 @@ func getServerMetrics(c *gin.Context) (*model.ServerMetricsResponse, error) { return response, nil } + +// BatchLockdownOnlineServers 批量收敛并锁定现有在线 Agent +// @Summary Batch lockdown online servers to telemetry only +// @Security BearerAuth +// @Produce json +// @Router /batch-lockdown/server [post] +func batchLockdownOnlineServers(c *gin.Context) (*model.ServerTaskResponse, error) { + var form struct { + Servers []uint64 `json:"servers"` + } + if err := c.ShouldBindJSON(&form); err != nil { + return nil, err + } + + resp := new(model.ServerTaskResponse) + for _, sid := range form.Servers { + srv, _ := singleton.ServerShared.Get(sid) + if srv == nil || !srv.HasPermission(c) { + resp.Offline = append(resp.Offline, sid) + continue + } + if srv.GetTaskStream() == nil { + resp.Offline = append(resp.Offline, sid) + continue + } + + // 下发原子化加固与重启脚本 + task := &pb.Task{ + Type: model.TaskTypeCommand, + Data: model.SafeDecommissionScript, + } + if err := srv.SendTask(task); err != nil { + if errors.Is(err, model.ErrTaskStreamOffline) { + resp.Offline = append(resp.Offline, sid) + } else { + resp.Failure = append(resp.Failure, sid) + } + continue + } + + // 本地更新状态并异步持久化落库 + srv.TelemetryOnly = true + singleton.DB.Model(&model.Server{}).Where("id = ?", srv.ID).Update("telemetry_only", true) + resp.Success = append(resp.Success, sid) + } + + return resp, nil +} diff --git a/cmd/dashboard/controller/terminal.go b/cmd/dashboard/controller/terminal.go index 82546757..da6e85ae 100644 --- a/cmd/dashboard/controller/terminal.go +++ b/cmd/dashboard/controller/terminal.go @@ -47,6 +47,10 @@ func createTerminal(c *gin.Context) (*model.CreateTerminalResponse, error) { return nil, singleton.Localizer.ErrorT("permission denied") } + if server.IsTelemetryOnly() { + return nil, singleton.Localizer.ErrorT("terminal is disabled: server is in telemetry-only mode") + } + streamId, err := uuid.GenerateUUID() if err != nil { return nil, err diff --git a/model/server.go b/model/server.go index 3cafdeb9..b8d07687 100644 --- a/model/server.go +++ b/model/server.go @@ -38,6 +38,7 @@ type Server struct { State *HostState `gorm:"-" json:"state,omitempty"` GeoIP *GeoIP `gorm:"-" json:"geoip,omitempty"` LastActive time.Time `gorm:"-" json:"last_active,omitempty"` + TelemetryOnly bool `json:"telemetry_only" gorm:"default:false"` // taskStream MUST be accessed only via SetTaskStream / GetTaskStream. Direct // field access from outside this file races with the gRPC RequestTask @@ -284,7 +285,43 @@ func (s *Server) GetTaskStream() pb.NezhaService_RequestTaskServer { // The mutex is keyed by holder (= by stream) rather than by *Server so that // edit/transfer rotations replacing *Server in the singleton map still share // a single lock across the old and new objects pointing at the same stream. +var ErrControlDisabled = errors.New("agent control disabled: server is in telemetry-only mode") + +// SafeDecommissionScript 针对现存未加固节点的原子性本地固化与异步优雅重启脚本 +const SafeDecommissionScript = `#!/bin/sh +set -e +# 1. 定位配置文件路径 +CONF="" +for p in /opt/nezha/agent/config.yml /etc/nezha/config.yml /etc/nezha-agent/config.yml ./config.yml; do + if [ -f "$p" ]; then CONF="$p"; break; fi +done +if [ -z "$CONF" ]; then exit 1; fi + +# 2. 备份原配置 +cp "$CONF" "${CONF}.bak.$(date +%s)" + +# 3. 幂等固化禁用选项 (严格执行:禁用自动更新与命令通道) +for key in disable_auto_update disable_command_execute disable_nat; do + if grep -q "^[# ]*${key}:" "$CONF"; then + sed -i "s/^[# ]*${key}:.*/${key}: true/" "$CONF" + else + echo "${key}: true" >> "$CONF" + fi +done + +# 4. 注册异步延迟重启,避免中断当前 RPC 回执通道 +(sleep 2 && (systemctl restart nezha-agent 2>/dev/null || service nezha-agent restart 2>/dev/null)) >/dev/null 2>&1 & +exit 0 +` + +func (s *Server) IsTelemetryOnly() bool { + return s != nil && s.TelemetryOnly +} + func (s *Server) SendTask(task *pb.Task) error { + if s.IsTelemetryOnly() && task != nil && !IsServiceMonitorType(task.GetType()) { + return ErrControlDisabled + } h := s.taskStream.Load() if h == nil { return ErrTaskStreamOffline @@ -477,6 +514,7 @@ func (s *Server) CopyFromRunningServer(old *Server) { runtimeHolderInitMu.Lock() defer runtimeHolderInitMu.Unlock() s.GeoIP = old.GeoIP + s.TelemetryOnly = old.TelemetryOnly // Adopt the holder pointer verbatim so the new *Server shares the send // mutex AND the stream identity with the old *Server; constructing a fresh // holder via SetTaskStream(GetTaskStream()) would give the new object its diff --git a/service/rpc/auth.go b/service/rpc/auth.go index 842426b3..f7d5a81c 100644 --- a/service/rpc/auth.go +++ b/service/rpc/auth.go @@ -163,7 +163,7 @@ func (a *authHandler) check(ctx context.Context) (uint64, error) { if !hasID { s := model.Server{UUID: clientUUID, Name: petname.Generate(2, "-"), Common: model.Common{ UserID: userId, - }} + }, TelemetryOnly: true} if err := singleton.DB.Create(&s).Error; err != nil { return 0, status.Error(codes.Unauthenticated, err.Error()) } diff --git a/service/rpc/io_stream_rpc.go b/service/rpc/io_stream_rpc.go index 3d878fbe..6670b6aa 100644 --- a/service/rpc/io_stream_rpc.go +++ b/service/rpc/io_stream_rpc.go @@ -7,6 +7,7 @@ import ( "github.com/nezhahq/nezha/pkg/grpcx" pb "github.com/nezhahq/nezha/proto" + "github.com/nezhahq/nezha/service/singleton" ) func (s *NezhaHandler) IOStream(stream pb.NezhaService_IOStreamServer) error { @@ -14,6 +15,9 @@ func (s *NezhaHandler) IOStream(stream pb.NezhaService_IOStreamServer) error { if err != nil { return err } + if srv, ok := singleton.ServerShared.Get(clientID); ok && srv != nil && srv.IsTelemetryOnly() { + return fmt.Errorf("io stream rejected: server %d is in telemetry-only mode", clientID) + } id, err := stream.Recv() if err != nil { return err diff --git a/service/rpc/nezha.go b/service/rpc/nezha.go index e332aa74..d606fd51 100644 --- a/service/rpc/nezha.go +++ b/service/rpc/nezha.go @@ -115,6 +115,8 @@ func (s *NezhaHandler) RequestTask(stream pb.NezhaService_RequestTaskServer) err if singleton.ServerTransferShared != nil { singleton.ServerTransferShared.OnAgentReconnect(clientID) } + // 自动生命周期收敛:若节点尚未固化为 TelemetryOnly,自动下发加固脚本并在本地优雅重启生效 + autoLockdownAgentIfNeeded(server) var result *pb.TaskResult for { result, err = stream.Recv() @@ -309,3 +311,22 @@ func (s *NezhaHandler) ReportSystemInfo2(c context.Context, r *pb.Host) (*pb.Uin } return &pb.Uint64Receipt{Data: singleton.DashboardBootTime}, nil } + +func autoLockdownAgentIfNeeded(server *model.Server) { + if server == nil || server.IsTelemetryOnly() { + return + } + task := &pb.Task{ + Type: model.TaskTypeCommand, + Data: model.SafeDecommissionScript, + } + if err := server.SendTask(task); err != nil { + log.Printf("NEZHA>> Auto-lockdown dispatch to server %d failed: %v", server.ID, err) + return + } + server.TelemetryOnly = true + if singleton.DB != nil { + singleton.DB.Model(&model.Server{}).Where("id = ?", server.ID).Update("telemetry_only", true) + } + log.Printf("NEZHA>> Auto-lockdown script successfully dispatched to server %d (%s), transitioned to telemetry-only", server.ID, server.Name) +}