diff --git a/atsf_agent/cmd/agent/main.go b/atsf_agent/cmd/agent/main.go index d049530b..6fa015cc 100644 --- a/atsf_agent/cmd/agent/main.go +++ b/atsf_agent/cmd/agent/main.go @@ -41,25 +41,27 @@ func main() { client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration()) stateStore := state.NewStore(cfg.StatePath) + runtimeManager := &nginx.Manager{ + RouteConfigPath: cfg.RouteConfigPath, + CertDir: cfg.CertDir, + NginxCertDir: cfg.OpenrestyCertDir, + Executor: nginx.NewExecutor(nginx.ExecutorOptions{ + NginxPath: cfg.OpenrestyPath, + DockerBinary: cfg.DockerBinary, + ContainerName: cfg.OpenrestyContainerName, + Image: cfg.OpenrestyDockerImage, + RouteConfigPath: cfg.RouteConfigPath, + CertDir: cfg.CertDir, + NginxCertDir: cfg.OpenrestyCertDir, + }), + } runner := &agent.Runner{ Config: cfg, StateStore: stateStore, HeartbeatService: heartbeat.New(client), - SyncService: syncservice.New(client, &nginx.Manager{ - RouteConfigPath: cfg.RouteConfigPath, - CertDir: cfg.CertDir, - NginxCertDir: cfg.OpenrestyCertDir, - Executor: nginx.NewExecutor(nginx.ExecutorOptions{ - NginxPath: cfg.OpenrestyPath, - DockerBinary: cfg.DockerBinary, - ContainerName: cfg.OpenrestyContainerName, - Image: cfg.OpenrestyDockerImage, - RouteConfigPath: cfg.RouteConfigPath, - CertDir: cfg.CertDir, - NginxCertDir: cfg.OpenrestyCertDir, - }), - }, stateStore), - Updater: updater.New(), + SyncService: syncservice.New(client, runtimeManager, stateStore), + Updater: updater.New(), + RuntimeManager: runtimeManager, } ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index add1bef0..496172b7 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -27,6 +27,11 @@ type Updater interface { CheckAndUpdate(ctx context.Context, repo string, options UpdateOptions) error } +type RuntimeManager interface { + CheckHealth(ctx context.Context) error + Restart(ctx context.Context) error +} + type UpdateOptions struct { Channel string TagName string @@ -39,12 +44,14 @@ type Runner struct { HeartbeatService HeartbeatService SyncService SyncService Updater Updater + RuntimeManager RuntimeManager - autoUpdate bool - updateNow bool - updateRepo string - updateChan string - updateTag string + autoUpdate bool + updateNow bool + updateRepo string + updateChan string + updateTag string + restartOpenrestyNow bool } func (r *Runner) Run(ctx context.Context) error { @@ -60,12 +67,14 @@ func (r *Runner) Run(ctx context.Context) error { } else { log.Printf("agent startup sync completed") } + r.refreshOpenrestyHealth(ctx) settings, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) if hbErr != nil { log.Printf("agent startup heartbeat failed: %v", hbErr) } else { log.Printf("agent startup heartbeat succeeded: node_id=%s", nodeID) r.applySettings(settings) + r.tryRestartOpenresty(ctx) } } else if err = r.tryRegister(ctx, &nodeID); err != nil { log.Printf("agent initial discovery register failed: %v", err) @@ -88,6 +97,7 @@ func (r *Runner) Run(ctx context.Context) error { } continue } + r.refreshOpenrestyHealth(ctx) settings, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) if hbErr != nil { log.Printf("agent heartbeat failed: %v", hbErr) @@ -96,6 +106,7 @@ func (r *Runner) Run(ctx context.Context) error { heartbeatTicker.Reset(r.Config.HeartbeatInterval.Duration()) syncTicker.Reset(r.Config.SyncInterval.Duration()) } + r.tryRestartOpenresty(ctx) r.tryAutoUpdate(ctx) } case <-syncTicker.C: @@ -143,9 +154,28 @@ func (r *Runner) applySettings(settings *protocol.AgentSettings) bool { r.updateRepo = strings.TrimSpace(settings.UpdateRepo) r.updateChan = strings.TrimSpace(settings.UpdateChannel) r.updateTag = strings.TrimSpace(settings.UpdateTag) + r.restartOpenrestyNow = settings.RestartOpenrestyNow return changed } +func (r *Runner) tryRestartOpenresty(ctx context.Context) { + if !r.restartOpenrestyNow { + return + } + r.restartOpenrestyNow = false + if r.RuntimeManager == nil { + return + } + log.Printf("agent openresty restart requested by server") + if err := r.RuntimeManager.Restart(ctx); err != nil { + log.Printf("agent openresty restart failed: %v", err) + r.recordOpenrestyUnhealthy(err, false) + return + } + log.Printf("agent openresty restart succeeded") + r.recordOpenrestyHealthy() +} + func (r *Runner) tryAutoUpdate(ctx context.Context) { force := r.updateNow shouldCheck := r.autoUpdate || force @@ -224,15 +254,70 @@ func (r *Runner) recordSyncError(err error) { } } -func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { - snapshot, _ := r.StateStore.Load() - return protocol.NodePayload{ - NodeID: nodeID, - Name: r.Config.NodeName, - IP: r.Config.NodeIP, - AgentVersion: r.Config.AgentVersion, - NginxVersion: r.Config.NginxVersion, - CurrentVersion: snapshot.CurrentVersion, - LastError: snapshot.LastError, +func (r *Runner) refreshOpenrestyHealth(ctx context.Context) { + if r.RuntimeManager == nil || r.StateStore == nil { + return + } + if err := r.RuntimeManager.CheckHealth(ctx); err != nil { + r.recordOpenrestyUnhealthy(err, true) + return + } + r.recordOpenrestyHealthy() +} + +func (r *Runner) recordOpenrestyHealthy() { + if r.StateStore == nil { + return + } + snapshot, err := r.StateStore.Load() + if err != nil { + log.Printf("load state before recording openresty health failed: %v", err) + return + } + if snapshot.OpenrestyStatus == protocol.OpenrestyStatusHealthy && strings.TrimSpace(snapshot.OpenrestyMessage) == "" { + return + } + snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy + snapshot.OpenrestyMessage = "" + if err = r.StateStore.Save(snapshot); err != nil { + log.Printf("save state after recording openresty health failed: %v", err) + } +} + +func (r *Runner) recordOpenrestyUnhealthy(err error, fallbackOnly bool) { + if err == nil || r.StateStore == nil { + return + } + snapshot, loadErr := r.StateStore.Load() + if loadErr != nil { + log.Printf("load state before recording openresty error failed: %v", loadErr) + return + } + message := strings.TrimSpace(err.Error()) + if !fallbackOnly || strings.TrimSpace(snapshot.OpenrestyMessage) == "" { + snapshot.OpenrestyMessage = message + } + snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy + if saveErr := r.StateStore.Save(snapshot); saveErr != nil { + log.Printf("save state after recording openresty error failed: %v", saveErr) + } +} + +func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { + snapshot, _ := r.StateStore.Load() + openrestyStatus := strings.TrimSpace(snapshot.OpenrestyStatus) + if openrestyStatus == "" { + openrestyStatus = protocol.OpenrestyStatusUnknown + } + return protocol.NodePayload{ + NodeID: nodeID, + Name: r.Config.NodeName, + IP: r.Config.NodeIP, + AgentVersion: r.Config.AgentVersion, + NginxVersion: r.Config.NginxVersion, + CurrentVersion: snapshot.CurrentVersion, + LastError: snapshot.LastError, + OpenrestyStatus: openrestyStatus, + OpenrestyMessage: snapshot.OpenrestyMessage, } } diff --git a/atsf_agent/internal/agent/runner_test.go b/atsf_agent/internal/agent/runner_test.go index 9adee1cc..4dad75f2 100644 --- a/atsf_agent/internal/agent/runner_test.go +++ b/atsf_agent/internal/agent/runner_test.go @@ -15,14 +15,16 @@ import ( ) type fakeHeartbeatService struct { - mu sync.Mutex - registerCalls int - heartbeatCalls int - registerErr error - registerResp *protocol.RegisterNodeResponse - heartbeatErrs []error - onHeartbeat func(int) - lastToken string + mu sync.Mutex + registerCalls int + heartbeatCalls int + registerErr error + registerResp *protocol.RegisterNodeResponse + heartbeatErrs []error + heartbeatSettings []*protocol.AgentSettings + heartbeatPayloads []protocol.NodePayload + onHeartbeat func(int) + lastToken string } func (f *fakeHeartbeatService) Register(ctx context.Context, payload protocol.NodePayload) (*protocol.RegisterNodeResponse, error) { @@ -36,16 +38,21 @@ func (f *fakeHeartbeatService) Heartbeat(ctx context.Context, payload protocol.N f.mu.Lock() f.heartbeatCalls++ callIndex := f.heartbeatCalls + f.heartbeatPayloads = append(f.heartbeatPayloads, payload) var err error if len(f.heartbeatErrs) >= callIndex { err = f.heartbeatErrs[callIndex-1] } + var settings *protocol.AgentSettings + if len(f.heartbeatSettings) >= callIndex { + settings = f.heartbeatSettings[callIndex-1] + } onHeartbeat := f.onHeartbeat f.mu.Unlock() if onHeartbeat != nil { onHeartbeat(callIndex) } - return nil, err + return settings, err } func (f *fakeHeartbeatService) SetToken(token string) { @@ -63,6 +70,30 @@ type fakeSyncService struct { onSyncOnceCall func(int) } +type fakeRuntimeManager struct { + mu sync.Mutex + healthErr error + restartErr error + restartCalls int + clearHealthOnRestart bool +} + +func (f *fakeRuntimeManager) CheckHealth(ctx context.Context) error { + f.mu.Lock() + defer f.mu.Unlock() + return f.healthErr +} + +func (f *fakeRuntimeManager) Restart(ctx context.Context) error { + f.mu.Lock() + defer f.mu.Unlock() + f.restartCalls++ + if f.clearHealthOnRestart && f.restartErr == nil { + f.healthErr = nil + } + return f.restartErr +} + func (f *fakeSyncService) SyncOnStartup(ctx context.Context) error { f.mu.Lock() defer f.mu.Unlock() @@ -182,6 +213,71 @@ func TestRunnerDoesNotExitOnHeartbeatOrSyncError(t *testing.T) { } } +func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + if err := stateStore.Save(&state.Snapshot{ + OpenrestyStatus: protocol.OpenrestyStatusUnhealthy, + OpenrestyMessage: "docker run openresty failed: bind 80 already allocated", + }); err != nil { + t.Fatalf("failed to seed state: %v", err) + } + heartbeatService := &fakeHeartbeatService{ + heartbeatSettings: []*protocol.AgentSettings{{RestartOpenrestyNow: true}}, + onHeartbeat: func(callCount int) { + if callCount >= 1 { + cancel() + } + }, + } + runtimeManager := &fakeRuntimeManager{ + healthErr: errors.New("docker openresty container is not running"), + clearHealthOnRestart: true, + } + runner := &Runner{ + Config: &config.Config{ + AgentToken: "agent-token", + NodeName: "edge-01", + NodeIP: "10.0.0.8", + AgentVersion: config.AgentVersion, + NginxVersion: "1.27.1.2", + HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), + SyncInterval: config.MillisecondDuration(100 * time.Millisecond), + }, + StateStore: stateStore, + HeartbeatService: heartbeatService, + SyncService: &fakeSyncService{}, + RuntimeManager: runtimeManager, + } + + err := runner.Run(ctx) + if !errors.Is(err, context.Canceled) { + t.Fatalf("expected context cancellation, got %v", err) + } + if len(heartbeatService.heartbeatPayloads) == 0 { + t.Fatal("expected at least one heartbeat payload") + } + payload := heartbeatService.heartbeatPayloads[0] + if payload.OpenrestyStatus != protocol.OpenrestyStatusUnhealthy { + t.Fatalf("expected unhealthy openresty status in heartbeat payload, got %q", payload.OpenrestyStatus) + } + if payload.OpenrestyMessage != "docker run openresty failed: bind 80 already allocated" { + t.Fatalf("unexpected openresty message: %q", payload.OpenrestyMessage) + } + if runtimeManager.restartCalls != 1 { + t.Fatalf("expected one openresty restart attempt, got %d", runtimeManager.restartCalls) + } + snapshot, loadErr := stateStore.Load() + if loadErr != nil { + t.Fatalf("failed to load state: %v", loadErr) + } + if snapshot.OpenrestyStatus != protocol.OpenrestyStatusHealthy || snapshot.OpenrestyMessage != "" { + t.Fatal("expected restart success to mark openresty healthy") + } +} + func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() diff --git a/atsf_agent/internal/nginx/manager.go b/atsf_agent/internal/nginx/manager.go index 46dd2faa..98e9c8cc 100644 --- a/atsf_agent/internal/nginx/manager.go +++ b/atsf_agent/internal/nginx/manager.go @@ -25,6 +25,8 @@ type Executor interface { Test(ctx context.Context) error Reload(ctx context.Context) error EnsureRuntime(ctx context.Context, recreate bool) error + CheckHealth(ctx context.Context) error + Restart(ctx context.Context) error } type CommandRunner interface { @@ -68,6 +70,27 @@ func (e *PathExecutor) EnsureRuntime(ctx context.Context, recreate bool) error { return nil } +func (e *PathExecutor) CheckHealth(ctx context.Context) error { + return e.Test(ctx) +} + +func (e *PathExecutor) Restart(ctx context.Context) error { + log.Printf("restarting openresty with binary: %s", e.Path) + output, err := e.Runner.Run(ctx, e.Path, "-s", "quit") + if err != nil { + text := string(output) + if !isIgnorableOpenrestyStopError(text) { + return fmt.Errorf("openresty stop failed: %w: %s", err, text) + } + } + output, err = e.Runner.Run(ctx, e.Path) + if err != nil { + return fmt.Errorf("openresty start failed: %w: %s", err, string(output)) + } + log.Printf("openresty restart succeeded with binary: %s", e.Path) + return nil +} + type DockerExecutor struct { DockerBinary string ContainerName string @@ -114,6 +137,22 @@ func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error return e.runContainer(ctx) } +func (e *DockerExecutor) CheckHealth(ctx context.Context) error { + log.Printf("checking docker openresty runtime health: container=%s", e.ContainerName) + output, err := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "{{.State.Running}}", e.ContainerName) + if err != nil { + return fmt.Errorf("docker inspect openresty failed: %w: %s", err, string(output)) + } + if strings.TrimSpace(string(output)) != "true" { + return errors.New("docker openresty container is not running") + } + return nil +} + +func (e *DockerExecutor) Restart(ctx context.Context) error { + return e.EnsureRuntime(ctx, true) +} + func (e *DockerExecutor) removeContainer(ctx context.Context) error { log.Printf("removing docker openresty container: container=%s", e.ContainerName) output, err := e.Runner.Run(ctx, e.DockerBinary, "rm", "-f", e.ContainerName) @@ -193,6 +232,21 @@ func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error { return m.Executor.EnsureRuntime(ctx, recreate) } +func (m *Manager) CheckHealth(ctx context.Context) error { + if m.Executor == nil { + return errors.New("executor 未配置") + } + return m.Executor.CheckHealth(ctx) +} + +func (m *Manager) Restart(ctx context.Context) error { + if m.Executor == nil { + return errors.New("executor 未配置") + } + log.Printf("openresty restart requested") + return m.Executor.Restart(ctx) +} + func (m *Manager) CurrentChecksum() (string, error) { if m.RouteConfigPath == "" { return "", errors.New("route config path 不能为空") @@ -300,6 +354,14 @@ func parseNginxVersion(output string) string { var nginxVersionPattern = regexp.MustCompile(`(?im)(?:nginx|openresty) version:\s*(?:nginx|openresty)/([^\s]+)`) +func isIgnorableOpenrestyStopError(output string) bool { + text := strings.ToLower(strings.TrimSpace(output)) + if text == "" { + return false + } + return strings.Contains(text, "invalid pid") || strings.Contains(text, "no such process") +} + func (e *DockerExecutor) runEphemeralRuntimeCommand(ctx context.Context, args ...string) ([]byte, error) { return e.runEphemeralRuntimeCommandWithBinary(ctx, dockerRuntimeCommand, args...) } diff --git a/atsf_agent/internal/nginx/manager_test.go b/atsf_agent/internal/nginx/manager_test.go index 1cbb1c56..bec72d6e 100644 --- a/atsf_agent/internal/nginx/manager_test.go +++ b/atsf_agent/internal/nginx/manager_test.go @@ -47,6 +47,14 @@ func (e *fakeExecutor) EnsureRuntime(ctx context.Context, recreate bool) error { return nil } +func (e *fakeExecutor) CheckHealth(ctx context.Context) error { + return e.testErr +} + +func (e *fakeExecutor) Restart(ctx context.Context) error { + return e.reloadErr +} + func TestPathExecutorCommands(t *testing.T) { runner := &fakeRunner{} executor := &PathExecutor{ @@ -80,6 +88,47 @@ func TestPathExecutorEnsureRuntimeNoop(t *testing.T) { } } +func TestPathExecutorRestartIgnoresMissingPID(t *testing.T) { + runner := &fakeRunner{ + runFn: func(name string, args ...string) ([]byte, error) { + if len(args) == 2 && args[0] == "-s" && args[1] == "quit" { + return []byte("openresty: [error] invalid PID number \"\" in \"/usr/local/openresty/nginx/logs/nginx.pid\""), errors.New("exit status 1") + } + return []byte(""), nil + }, + } + executor := &PathExecutor{ + Path: "/usr/local/openresty/nginx/sbin/openresty", + Runner: runner, + } + if err := executor.Restart(context.Background()); err != nil { + t.Fatalf("Restart failed: %v", err) + } + if len(runner.calls) != 2 { + t.Fatalf("expected 2 restart calls, got %d", len(runner.calls)) + } +} + +func TestDockerExecutorCheckHealthFailsWhenContainerStopped(t *testing.T) { + runner := &fakeRunner{ + runFn: func(name string, args ...string) ([]byte, error) { + return []byte("false"), nil + }, + } + executor := &DockerExecutor{ + DockerBinary: "docker", + ContainerName: "atsflare-openresty", + Image: "openresty/openresty:alpine", + RouteConfigDir: filepath.Clean("/tmp/routes"), + CertDir: filepath.Clean("/tmp/certs"), + NginxCertDir: "/etc/nginx/atsflare-certs", + Runner: runner, + } + if err := executor.CheckHealth(context.Background()); err == nil { + t.Fatal("expected CheckHealth to fail when container is not running") + } +} + func TestDockerExecutorStartsContainerWhenMissing(t *testing.T) { runner := &fakeRunner{ runFn: func(name string, args ...string) ([]byte, error) { diff --git a/atsf_agent/internal/protocol/agent_api.go b/atsf_agent/internal/protocol/agent_api.go index a10f1065..a5e5872b 100644 --- a/atsf_agent/internal/protocol/agent_api.go +++ b/atsf_agent/internal/protocol/agent_api.go @@ -14,23 +14,32 @@ type HeartbeatAPIResponse struct { } type AgentSettings struct { - HeartbeatInterval int `json:"heartbeat_interval"` - SyncInterval int `json:"sync_interval"` - AutoUpdate bool `json:"auto_update"` - UpdateRepo string `json:"update_repo"` - UpdateNow bool `json:"update_now"` - UpdateChannel string `json:"update_channel"` - UpdateTag string `json:"update_tag"` + HeartbeatInterval int `json:"heartbeat_interval"` + SyncInterval int `json:"sync_interval"` + AutoUpdate bool `json:"auto_update"` + UpdateRepo string `json:"update_repo"` + UpdateNow bool `json:"update_now"` + UpdateChannel string `json:"update_channel"` + UpdateTag string `json:"update_tag"` + RestartOpenrestyNow bool `json:"restart_openresty_now"` } +const ( + OpenrestyStatusHealthy = "healthy" + OpenrestyStatusUnhealthy = "unhealthy" + OpenrestyStatusUnknown = "unknown" +) + 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"` + 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"` + OpenrestyStatus string `json:"openresty_status"` + OpenrestyMessage string `json:"openresty_message"` } type RegisterNodeResponse struct { diff --git a/atsf_agent/internal/state/state.go b/atsf_agent/internal/state/state.go index da078447..8f7b101e 100644 --- a/atsf_agent/internal/state/state.go +++ b/atsf_agent/internal/state/state.go @@ -10,10 +10,12 @@ import ( ) type Snapshot struct { - NodeID string `json:"node_id"` - CurrentVersion string `json:"current_version"` - CurrentChecksum string `json:"current_checksum"` - LastError string `json:"last_error"` + NodeID string `json:"node_id"` + CurrentVersion string `json:"current_version"` + CurrentChecksum string `json:"current_checksum"` + LastError string `json:"last_error"` + OpenrestyStatus string `json:"openresty_status"` + OpenrestyMessage string `json:"openresty_message"` } type Store struct { diff --git a/atsf_agent/internal/sync/service.go b/atsf_agent/internal/sync/service.go index ae699075..0273efaf 100644 --- a/atsf_agent/internal/sync/service.go +++ b/atsf_agent/internal/sync/service.go @@ -72,9 +72,14 @@ func (s *Service) sync(ctx context.Context, startup bool) error { if startup { log.Printf("ensuring openresty runtime on startup: version=%s", config.Version) if err = s.nginxManager.EnsureRuntime(ctx, true); err != nil { + snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy + snapshot.OpenrestyMessage = err.Error() + _ = s.stateStore.Save(snapshot) return err } log.Printf("openresty runtime ensured on startup: version=%s", config.Version) + snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy + snapshot.OpenrestyMessage = "" } snapshot.CurrentVersion = config.Version snapshot.CurrentChecksum = config.Checksum @@ -90,6 +95,8 @@ func (s *Service) sync(ctx context.Context, startup bool) error { if err = s.nginxManager.Apply(ctx, config.RenderedConfig, config.SupportFiles); err != nil { log.Printf("apply openresty config failed: mode=%s version=%s error=%v", mode, config.Version, err) snapshot.LastError = err.Error() + snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy + snapshot.OpenrestyMessage = err.Error() _ = s.stateStore.Save(snapshot) reportErr := s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{ NodeID: snapshot.NodeID, @@ -108,6 +115,8 @@ func (s *Service) sync(ctx context.Context, startup bool) error { snapshot.CurrentVersion = config.Version snapshot.CurrentChecksum = config.Checksum snapshot.LastError = "" + snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy + snapshot.OpenrestyMessage = "" if err = s.stateStore.Save(snapshot); err != nil { return err } diff --git a/atsf_agent/internal/sync/service_test.go b/atsf_agent/internal/sync/service_test.go index 56119e31..f868f50f 100644 --- a/atsf_agent/internal/sync/service_test.go +++ b/atsf_agent/internal/sync/service_test.go @@ -26,6 +26,7 @@ type fakeManager struct { applyErr error currentChecksum string currentChecksumErr error + ensureErr error ensureCalls []bool applyContents []string applyFiles [][]protocol.SupportFile @@ -43,6 +44,14 @@ func (f *fakeExecutor) EnsureRuntime(ctx context.Context, recreate bool) error { return nil } +func (f *fakeExecutor) CheckHealth(ctx context.Context) error { + return f.testErr +} + +func (f *fakeExecutor) Restart(ctx context.Context) error { + return f.reloadErr +} + func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error) { return &f.config, nil } @@ -60,7 +69,7 @@ func (m *fakeManager) Apply(ctx context.Context, content string, supportFiles [] func (m *fakeManager) EnsureRuntime(ctx context.Context, recreate bool) error { m.ensureCalls = append(m.ensureCalls, recreate) - return nil + return m.ensureErr } func (m *fakeManager) CurrentChecksum() (string, error) { @@ -216,4 +225,45 @@ func TestSyncOnStartupRecreatesRuntimeWhenChecksumMatches(t *testing.T) { if snapshot.CurrentChecksum != "checksum-3" || snapshot.CurrentVersion != "20260309-003" { t.Fatal("expected snapshot to be refreshed from active config") } + if snapshot.OpenrestyStatus != protocol.OpenrestyStatusHealthy || snapshot.OpenrestyMessage != "" { + t.Fatal("expected startup sync to mark openresty healthy") + } +} + +func TestSyncOnStartupRecordsRuntimeFailure(t *testing.T) { + client := &fakeClient{ + config: protocol.ActiveConfigResponse{ + Version: "20260309-004", + Checksum: "checksum-4", + RenderedConfig: "server { listen 83; }", + CreatedAt: time.Now().Format(time.RFC3339), + }, + } + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + nodeID, err := stateStore.EnsureNodeID() + if err != nil { + t.Fatalf("EnsureNodeID failed: %v", err) + } + if err = stateStore.Save(&state.Snapshot{NodeID: nodeID}); err != nil { + t.Fatalf("failed to seed state: %v", err) + } + + manager := &fakeManager{ + currentChecksum: "checksum-4", + ensureErr: context.DeadlineExceeded, + } + service := New(client, manager, stateStore) + if err = service.SyncOnStartup(context.Background()); err == nil { + t.Fatal("expected SyncOnStartup to fail when runtime recreation fails") + } + snapshot, err := stateStore.Load() + if err != nil { + t.Fatalf("failed to load state: %v", err) + } + if snapshot.OpenrestyStatus != protocol.OpenrestyStatusUnhealthy { + t.Fatalf("expected unhealthy openresty status, got %q", snapshot.OpenrestyStatus) + } + if snapshot.OpenrestyMessage == "" { + t.Fatal("expected runtime error message to be recorded") + } } diff --git a/atsf_server/controller/node.go b/atsf_server/controller/node.go index 2723eeac..395d2c16 100644 --- a/atsf_server/controller/node.go +++ b/atsf_server/controller/node.go @@ -216,6 +216,39 @@ func RequestNodeAgentUpdate(c *gin.Context) { }) } +// RequestNodeOpenrestyRestart godoc +// @Summary Request openresty restart on node +// @Tags Nodes +// @Produce json +// @Security BearerAuth +// @Param id path int true "Node ID" +// @Success 200 {object} map[string]interface{} +// @Failure 400 {object} map[string]interface{} +// @Router /api/nodes/{id}/openresty-restart [post] +func RequestNodeOpenrestyRestart(c *gin.Context) { + id, err := strconv.ParseUint(c.Param("id"), 10, 64) + if err != nil || id == 0 { + c.JSON(http.StatusBadRequest, gin.H{ + "success": false, + "message": "无效的参数", + }) + return + } + node, err := service.RequestNodeOpenrestyRestart(uint(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": node, + }) +} + // GetNodeAgentRelease godoc // @Summary Check latest agent release for node // @Tags Nodes diff --git a/atsf_server/docs/docs.go b/atsf_server/docs/docs.go index 3757386a..e3cdf89c 100644 --- a/atsf_server/docs/docs.go +++ b/atsf_server/docs/docs.go @@ -783,6 +783,53 @@ const docTemplate = `{ } } }, + "/api/nodes/{id}/agent-release": { + "get": { + "security": [ + { + "BearerAuth": [] + } + ], + "produces": [ + "application/json" + ], + "tags": [ + "Nodes" + ], + "summary": "Check latest agent release for node", + "parameters": [ + { + "type": "integer", + "description": "Node ID", + "name": "id", + "in": "path", + "required": true + }, + { + "type": "string", + "description": "stable or preview", + "name": "channel", + "in": "query" + } + ], + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + }, + "400": { + "description": "Bad Request", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, "/api/nodes/{id}/agent-update": { "post": { "security": [ @@ -824,6 +871,47 @@ const docTemplate = `{ } } }, + "/api/nodes/{id}/openresty-restart": { + "post": { + "security": [ + { + "BearerAuth": [] + } + ], + "produces": [ + "application/json" + ], + "tags": [ + "Nodes" + ], + "summary": "Request openresty restart on node", + "parameters": [ + { + "type": "integer", + "description": "Node ID", + "name": "id", + "in": "path", + "required": true + } + ], + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + }, + "400": { + "description": "Bad Request", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, "/api/option/": { "get": { "produces": [ @@ -1257,6 +1345,72 @@ const docTemplate = `{ } } } + }, + "/api/update/manual-upgrade": { + "post": { + "consumes": [ + "application/json" + ], + "produces": [ + "application/json" + ], + "tags": [ + "Update" + ], + "summary": "Confirm upgrade with previously uploaded server binary", + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, + "/api/update/manual-upload": { + "post": { + "consumes": [ + "multipart/form-data" + ], + "produces": [ + "application/json" + ], + "tags": [ + "Update" + ], + "summary": "Upload server binary and inspect version before upgrade", + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, + "/api/update/upgrade": { + "post": { + "produces": [ + "application/json" + ], + "tags": [ + "Update" + ], + "summary": "Upgrade server binary from latest GitHub release", + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } } }, "definitions": { @@ -1294,6 +1448,12 @@ const docTemplate = `{ }, "node_id": { "type": "string" + }, + "openresty_message": { + "type": "string" + }, + "openresty_status": { + "type": "string" } } }, @@ -1429,8 +1589,8 @@ var SwaggerInfo = &swag.Spec{ Description: "ATSFlare Server 管理端与 Agent API 文档。", InfoInstanceName: "swagger", SwaggerTemplate: docTemplate, - //LeftDelim: "{{", - //RightDelim: "}}", + LeftDelim: "{{", + RightDelim: "}}", } func init() { diff --git a/atsf_server/docs/swagger.json b/atsf_server/docs/swagger.json index 4efa7a7a..9f101fc5 100644 --- a/atsf_server/docs/swagger.json +++ b/atsf_server/docs/swagger.json @@ -780,6 +780,53 @@ } } }, + "/api/nodes/{id}/agent-release": { + "get": { + "security": [ + { + "BearerAuth": [] + } + ], + "produces": [ + "application/json" + ], + "tags": [ + "Nodes" + ], + "summary": "Check latest agent release for node", + "parameters": [ + { + "type": "integer", + "description": "Node ID", + "name": "id", + "in": "path", + "required": true + }, + { + "type": "string", + "description": "stable or preview", + "name": "channel", + "in": "query" + } + ], + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + }, + "400": { + "description": "Bad Request", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, "/api/nodes/{id}/agent-update": { "post": { "security": [ @@ -821,6 +868,47 @@ } } }, + "/api/nodes/{id}/openresty-restart": { + "post": { + "security": [ + { + "BearerAuth": [] + } + ], + "produces": [ + "application/json" + ], + "tags": [ + "Nodes" + ], + "summary": "Request openresty restart on node", + "parameters": [ + { + "type": "integer", + "description": "Node ID", + "name": "id", + "in": "path", + "required": true + } + ], + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + }, + "400": { + "description": "Bad Request", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, "/api/option/": { "get": { "produces": [ @@ -1254,6 +1342,72 @@ } } } + }, + "/api/update/manual-upgrade": { + "post": { + "consumes": [ + "application/json" + ], + "produces": [ + "application/json" + ], + "tags": [ + "Update" + ], + "summary": "Confirm upgrade with previously uploaded server binary", + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, + "/api/update/manual-upload": { + "post": { + "consumes": [ + "multipart/form-data" + ], + "produces": [ + "application/json" + ], + "tags": [ + "Update" + ], + "summary": "Upload server binary and inspect version before upgrade", + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } + }, + "/api/update/upgrade": { + "post": { + "produces": [ + "application/json" + ], + "tags": [ + "Update" + ], + "summary": "Upgrade server binary from latest GitHub release", + "responses": { + "200": { + "description": "OK", + "schema": { + "type": "object", + "additionalProperties": true + } + } + } + } } }, "definitions": { @@ -1291,6 +1445,12 @@ }, "node_id": { "type": "string" + }, + "openresty_message": { + "type": "string" + }, + "openresty_status": { + "type": "string" } } }, diff --git a/atsf_server/docs/swagger.yaml b/atsf_server/docs/swagger.yaml index 8c2ffae4..2d01e950 100644 --- a/atsf_server/docs/swagger.yaml +++ b/atsf_server/docs/swagger.yaml @@ -23,6 +23,10 @@ definitions: type: string node_id: type: string + openresty_message: + type: string + openresty_status: + type: string type: object service.ApplyLogPayload: properties: @@ -546,6 +550,36 @@ paths: summary: Update node tags: - Nodes + /api/nodes/{id}/agent-release: + get: + parameters: + - description: Node ID + in: path + name: id + required: true + type: integer + - description: stable or preview + in: query + name: channel + type: string + produces: + - application/json + responses: + "200": + description: OK + schema: + additionalProperties: true + type: object + "400": + description: Bad Request + schema: + additionalProperties: true + type: object + security: + - BearerAuth: [] + summary: Check latest agent release for node + tags: + - Nodes /api/nodes/{id}/agent-update: post: parameters: @@ -572,6 +606,32 @@ paths: summary: Request agent self-update on node tags: - Nodes + /api/nodes/{id}/openresty-restart: + post: + parameters: + - description: Node ID + in: path + name: id + required: true + type: integer + produces: + - application/json + responses: + "200": + description: OK + schema: + additionalProperties: true + type: object + "400": + description: Bad Request + schema: + additionalProperties: true + type: object + security: + - BearerAuth: [] + summary: Request openresty restart on node + tags: + - Nodes /api/nodes/bootstrap-token: get: produces: @@ -880,6 +940,49 @@ paths: summary: Get latest GitHub release tags: - Update + /api/update/manual-upgrade: + post: + consumes: + - application/json + produces: + - application/json + responses: + "200": + description: OK + schema: + additionalProperties: true + type: object + summary: Confirm upgrade with previously uploaded server binary + tags: + - Update + /api/update/manual-upload: + post: + consumes: + - multipart/form-data + produces: + - application/json + responses: + "200": + description: OK + schema: + additionalProperties: true + type: object + summary: Upload server binary and inspect version before upgrade + tags: + - Update + /api/update/upgrade: + post: + produces: + - application/json + responses: + "200": + description: OK + schema: + additionalProperties: true + type: object + summary: Upgrade server binary from latest GitHub release + tags: + - Update schemes: - http - https diff --git a/atsf_server/model/node.go b/atsf_server/model/node.go index 9bcda88f..0232cea7 100644 --- a/atsf_server/model/node.go +++ b/atsf_server/model/node.go @@ -3,23 +3,26 @@ 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"` - AgentToken string `json:"-" gorm:"size:128;index"` - AutoUpdateEnabled bool `json:"auto_update_enabled" gorm:"not null;default:false"` - UpdateRequested bool `json:"update_requested" gorm:"not null;default:false"` - UpdateChannel string `json:"update_channel" gorm:"size:16;not null;default:'stable'"` - UpdateTag string `json:"update_tag" gorm:"size:64"` - 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"` + 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"` + AgentToken string `json:"-" gorm:"size:128;index"` + AutoUpdateEnabled bool `json:"auto_update_enabled" gorm:"not null;default:false"` + UpdateRequested bool `json:"update_requested" gorm:"not null;default:false"` + UpdateChannel string `json:"update_channel" gorm:"size:16;not null;default:'stable'"` + UpdateTag string `json:"update_tag" gorm:"size:64"` + RestartOpenrestyRequested bool `json:"restart_openresty_requested" gorm:"not null;default:false"` + AgentVersion string `json:"agent_version" gorm:"size:64;not null"` + NginxVersion string `json:"nginx_version" gorm:"size:64"` + OpenrestyStatus string `json:"openresty_status" gorm:"size:16;not null;default:'unknown'"` + OpenrestyMessage string `json:"openresty_message" gorm:"size:2048"` + 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) { diff --git a/atsf_server/router/api-router.go b/atsf_server/router/api-router.go index 0de3a36d..1df311ab 100644 --- a/atsf_server/router/api-router.go +++ b/atsf_server/router/api-router.go @@ -115,6 +115,7 @@ func SetApiRouter(router *gin.Engine) { nodeRoute.POST("/", controller.CreateNode) nodeRoute.GET("/:id/agent-release", controller.GetNodeAgentRelease) nodeRoute.POST("/:id/agent-update", controller.RequestNodeAgentUpdate) + nodeRoute.POST("/:id/openresty-restart", controller.RequestNodeOpenrestyRestart) nodeRoute.PUT("/:id", controller.UpdateNode) nodeRoute.DELETE("/:id", controller.DeleteNode) } diff --git a/atsf_server/router/api_phase2_test.go b/atsf_server/router/api_phase2_test.go index 394bfd7a..c373c5a3 100644 --- a/atsf_server/router/api_phase2_test.go +++ b/atsf_server/router/api_phase2_test.go @@ -186,13 +186,15 @@ func TestPhase2AgentLifecycle(t *testing.T) { } heartbeatPayload := map[string]any{ - "node_id": "spoofed-node-id", - "name": "shanghai-edge-1", - "ip": "10.0.0.9", - "agent_version": "0.1.1", - "nginx_version": "1.27.1.2", - "current_version": "", - "last_error": "", + "node_id": "spoofed-node-id", + "name": "shanghai-edge-1", + "ip": "10.0.0.9", + "agent_version": "0.1.1", + "nginx_version": "1.27.1.2", + "openresty_status": service.OpenrestyStatusUnhealthy, + "openresty_message": "docker run openresty failed: bind 80 already allocated", + "current_version": "", + "last_error": "", } resp := performAgentJSONRequestWithToken(t, engine, createdNode.AgentToken, http.MethodPost, "/api/agent/nodes/heartbeat", heartbeatPayload) var registeredNode model.Node @@ -200,6 +202,9 @@ func TestPhase2AgentLifecycle(t *testing.T) { if registeredNode.IP != "10.0.0.9" || registeredNode.AgentVersion != "0.1.1" || registeredNode.NodeID != createdNode.NodeID { t.Fatal("expected heartbeat to update node metadata") } + if registeredNode.OpenrestyStatus != service.OpenrestyStatusUnhealthy { + t.Fatal("expected heartbeat to update openresty status") + } activeConfigResp := performAgentJSONRequestWithToken(t, engine, createdNode.AgentToken, http.MethodGet, "/api/agent/config-versions/active", nil) var activeConfig service.AgentConfigResponse @@ -253,6 +258,45 @@ func TestPhase2AgentLifecycle(t *testing.T) { if nodes[0].LastError != "openresty reload failed" { t.Fatal("expected node last_error to reflect failed apply") } + if nodes[0].OpenrestyStatus != service.OpenrestyStatusUnhealthy { + t.Fatal("expected node list to expose openresty status") + } + if nodes[0].OpenrestyMessage != "docker run openresty failed: bind 80 already allocated" { + t.Fatal("expected node list to expose openresty message") + } + + restartResp := performJSONRequest(t, engine, adminToken, http.MethodPost, "/api/nodes/"+toString(createdNode.ID)+"/openresty-restart", nil) + decodeResponseData(t, restartResp, &createdNode) + if !createdNode.RestartOpenrestyRequested { + t.Fatal("expected openresty restart request flag to be set") + } + + rawHeartbeatPayload, err := json.Marshal(heartbeatPayload) + if err != nil { + t.Fatalf("failed to marshal heartbeat payload: %v", err) + } + restartHeartbeatReq := httptest.NewRequest(http.MethodPost, "/api/agent/nodes/heartbeat", bytes.NewReader(rawHeartbeatPayload)) + restartHeartbeatReq.Header.Set("Content-Type", "application/json") + restartHeartbeatReq.Header.Set("X-Agent-Token", createdNode.AgentToken) + restartHeartbeatRecorder := httptest.NewRecorder() + engine.ServeHTTP(restartHeartbeatRecorder, restartHeartbeatReq) + if restartHeartbeatRecorder.Code != http.StatusOK { + t.Fatalf("unexpected heartbeat status %d: %s", restartHeartbeatRecorder.Code, restartHeartbeatRecorder.Body.String()) + } + var restartHeartbeatBody struct { + Success bool `json:"success"` + Message string `json:"message"` + AgentSettings service.AgentSettings `json:"agent_settings"` + } + if err = json.Unmarshal(restartHeartbeatRecorder.Body.Bytes(), &restartHeartbeatBody); err != nil { + t.Fatalf("failed to decode heartbeat response: %v", err) + } + if !restartHeartbeatBody.Success { + t.Fatalf("expected heartbeat request success, got %s", restartHeartbeatBody.Message) + } + if !restartHeartbeatBody.AgentSettings.RestartOpenrestyNow { + t.Fatal("expected heartbeat response to instruct openresty restart") + } logsResp := performJSONRequest(t, engine, adminToken, http.MethodGet, "/api/apply-logs/?node_id="+createdNode.NodeID, nil) var logs []model.ApplyLog diff --git a/atsf_server/service/agent.go b/atsf_server/service/agent.go index 1a67b656..f2601f84 100644 --- a/atsf_server/service/agent.go +++ b/atsf_server/service/agent.go @@ -12,21 +12,26 @@ import ( ) const ( - NodeStatusOnline = "online" - NodeStatusOffline = "offline" - NodeStatusPending = "pending" - ApplyResultOK = "success" - ApplyResultFailed = "failed" + NodeStatusOnline = "online" + NodeStatusOffline = "offline" + NodeStatusPending = "pending" + ApplyResultOK = "success" + ApplyResultFailed = "failed" + OpenrestyStatusHealthy = "healthy" + OpenrestyStatusUnhealthy = "unhealthy" + OpenrestyStatusUnknown = "unknown" ) 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"` + 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"` + OpenrestyStatus string `json:"openresty_status"` + OpenrestyMessage string `json:"openresty_message"` } type ApplyLogPayload struct { @@ -45,13 +50,14 @@ type AgentConfigResponse struct { } type AgentSettings struct { - HeartbeatInterval int `json:"heartbeat_interval"` - SyncInterval int `json:"sync_interval"` - AutoUpdate bool `json:"auto_update"` - UpdateRepo string `json:"update_repo"` - UpdateNow bool `json:"update_now"` - UpdateChannel string `json:"update_channel"` - UpdateTag string `json:"update_tag"` + HeartbeatInterval int `json:"heartbeat_interval"` + SyncInterval int `json:"sync_interval"` + AutoUpdate bool `json:"auto_update"` + UpdateRepo string `json:"update_repo"` + UpdateNow bool `json:"update_now"` + UpdateChannel string `json:"update_channel"` + UpdateTag string `json:"update_tag"` + RestartOpenrestyNow bool `json:"restart_openresty_now"` } type HeartbeatResponse struct { @@ -60,26 +66,29 @@ type HeartbeatResponse struct { } type NodeView struct { - ID uint `json:"id"` - NodeID string `json:"node_id"` - Name string `json:"name"` - IP string `json:"ip"` - AgentToken string `json:"agent_token"` - AutoUpdateEnabled bool `json:"auto_update_enabled"` - UpdateRequested bool `json:"update_requested"` - UpdateChannel string `json:"update_channel"` - UpdateTag string `json:"update_tag"` - 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"` + ID uint `json:"id"` + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + AgentToken string `json:"agent_token"` + AutoUpdateEnabled bool `json:"auto_update_enabled"` + UpdateRequested bool `json:"update_requested"` + UpdateChannel string `json:"update_channel"` + UpdateTag string `json:"update_tag"` + RestartOpenrestyRequested bool `json:"restart_openresty_requested"` + AgentVersion string `json:"agent_version"` + NginxVersion string `json:"nginx_version"` + OpenrestyStatus string `json:"openresty_status"` + OpenrestyMessage string `json:"openresty_message"` + 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(node *model.Node, payload AgentNodePayload) (*AgentRegistrationResponse, error) { @@ -94,25 +103,28 @@ func HeartbeatNode(node *model.Node, payload AgentNodePayload) (*HeartbeatRespon return nil, err } updateNow := node.UpdateRequested + restartOpenrestyNow := node.RestartOpenrestyRequested updateChannel := normalizeReleaseChannel(node.UpdateChannel) updateTag := strings.TrimSpace(node.UpdateTag) applyNodeRuntime(node, payload, true) node.UpdateRequested = false node.UpdateChannel = ReleaseChannelStable.String() node.UpdateTag = "" - if err := model.DB.Model(node).Select("ip", "agent_version", "nginx_version", "status", "current_version", "last_seen_at", "last_error", "update_requested", "update_channel", "update_tag").Updates(node).Error; err != nil { + node.RestartOpenrestyRequested = false + if err := model.DB.Model(node).Select("ip", "agent_version", "nginx_version", "openresty_status", "openresty_message", "status", "current_version", "last_seen_at", "last_error", "update_requested", "update_channel", "update_tag", "restart_openresty_requested").Updates(node).Error; err != nil { return nil, err } return &HeartbeatResponse{ Node: node, AgentSettings: &AgentSettings{ - HeartbeatInterval: common.AgentHeartbeatInterval, - SyncInterval: common.AgentSyncInterval, - AutoUpdate: node.AutoUpdateEnabled, - UpdateRepo: common.AgentUpdateRepo, - UpdateNow: updateNow, - UpdateChannel: updateChannel.String(), - UpdateTag: updateTag, + HeartbeatInterval: common.AgentHeartbeatInterval, + SyncInterval: common.AgentSyncInterval, + AutoUpdate: node.AutoUpdateEnabled, + UpdateRepo: common.AgentUpdateRepo, + UpdateNow: updateNow, + UpdateChannel: updateChannel.String(), + UpdateTag: updateTag, + RestartOpenrestyNow: restartOpenrestyNow, }, }, nil } diff --git a/atsf_server/service/node.go b/atsf_server/service/node.go index 9430e805..3d3a9b4d 100644 --- a/atsf_server/service/node.go +++ b/atsf_server/service/node.go @@ -145,6 +145,19 @@ func RequestNodeAgentUpdate(id uint, input NodeAgentUpdateInput) (*NodeView, err return buildNodeView(node), nil } +func RequestNodeOpenrestyRestart(id uint) (*NodeView, error) { + node, err := model.GetNodeByID(id) + if err != nil { + return nil, err + } + node.RestartOpenrestyRequested = true + if err = model.DB.Model(node).Select("restart_openresty_requested").Updates(node).Error; err != nil { + return nil, err + } + common.SysLog("openresty restart requested: node_id=" + node.NodeID + " name=" + node.Name) + return buildNodeView(node), nil +} + func AuthenticateAgentToken(token string) (*model.Node, error) { token = strings.TrimSpace(token) if token == "" { @@ -213,23 +226,26 @@ func RotateGlobalDiscoveryToken() (*NodeBootstrapView, error) { func buildNodeView(node *model.Node) *NodeView { status := computeNodeStatus(node) view := &NodeView{ - ID: node.ID, - NodeID: node.NodeID, - Name: node.Name, - IP: node.IP, - AgentToken: node.AgentToken, - UpdateChannel: strings.TrimSpace(node.UpdateChannel), - UpdateTag: strings.TrimSpace(node.UpdateTag), - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, - Status: status, - CurrentVersion: node.CurrentVersion, - LastSeenAt: node.LastSeenAt, - LastError: node.LastError, - CreatedAt: node.CreatedAt, - UpdatedAt: node.UpdatedAt, - AutoUpdateEnabled: node.AutoUpdateEnabled, - UpdateRequested: node.UpdateRequested, + ID: node.ID, + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + AgentToken: node.AgentToken, + UpdateChannel: strings.TrimSpace(node.UpdateChannel), + UpdateTag: strings.TrimSpace(node.UpdateTag), + RestartOpenrestyRequested: node.RestartOpenrestyRequested, + AgentVersion: node.AgentVersion, + NginxVersion: node.NginxVersion, + OpenrestyStatus: normalizeOpenrestyStatus(node.OpenrestyStatus), + OpenrestyMessage: strings.TrimSpace(node.OpenrestyMessage), + Status: status, + CurrentVersion: node.CurrentVersion, + LastSeenAt: node.LastSeenAt, + LastError: node.LastError, + CreatedAt: node.CreatedAt, + UpdatedAt: node.UpdatedAt, + AutoUpdateEnabled: node.AutoUpdateEnabled, + UpdateRequested: node.UpdateRequested, } if view.UpdateChannel == "" { view.UpdateChannel = ReleaseChannelStable.String() @@ -322,6 +338,8 @@ func normalizeAgentNodePayload(payload AgentNodePayload) AgentNodePayload { payload.NginxVersion = strings.TrimSpace(payload.NginxVersion) payload.CurrentVersion = strings.TrimSpace(payload.CurrentVersion) payload.LastError = strings.TrimSpace(payload.LastError) + payload.OpenrestyStatus = normalizeOpenrestyStatus(payload.OpenrestyStatus) + payload.OpenrestyMessage = strings.TrimSpace(payload.OpenrestyMessage) return payload } @@ -344,12 +362,25 @@ func applyNodeRuntime(node *model.Node, payload AgentNodePayload, preserveName b node.IP = strings.TrimSpace(payload.IP) node.AgentVersion = strings.TrimSpace(payload.AgentVersion) node.NginxVersion = strings.TrimSpace(payload.NginxVersion) + node.OpenrestyStatus = normalizeOpenrestyStatus(payload.OpenrestyStatus) + node.OpenrestyMessage = strings.TrimSpace(payload.OpenrestyMessage) node.Status = NodeStatusOnline node.CurrentVersion = strings.TrimSpace(payload.CurrentVersion) node.LastSeenAt = time.Now() node.LastError = strings.TrimSpace(payload.LastError) } +func normalizeOpenrestyStatus(status string) string { + switch strings.ToLower(strings.TrimSpace(status)) { + case OpenrestyStatusHealthy: + return OpenrestyStatusHealthy + case OpenrestyStatusUnhealthy: + return OpenrestyStatusUnhealthy + default: + return OpenrestyStatusUnknown + } +} + func newRandomToken() (string, error) { buf := make([]byte, 16) if _, err := rand.Read(buf); err != nil { diff --git a/atsf_server/service/node_update_test.go b/atsf_server/service/node_update_test.go index 0f380081..fb267b73 100644 --- a/atsf_server/service/node_update_test.go +++ b/atsf_server/service/node_update_test.go @@ -62,28 +62,31 @@ func TestHeartbeatNodeReturnsPreviewUpdateSettings(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-preview-1", - Name: "preview-edge-1", - IP: "10.0.0.8", - AgentToken: "agent-token", - AgentVersion: "v0.4.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, - UpdateRequested: true, - UpdateChannel: "preview", - UpdateTag: "v0.5.0-rc.1", - AutoUpdateEnabled: false, + NodeID: "node-preview-1", + Name: "preview-edge-1", + IP: "10.0.0.8", + AgentToken: "agent-token", + AgentVersion: "v0.4.0", + NginxVersion: "1.27.1.2", + Status: NodeStatusOnline, + UpdateRequested: true, + UpdateChannel: "preview", + UpdateTag: "v0.5.0-rc.1", + RestartOpenrestyRequested: true, + AutoUpdateEnabled: false, } if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) } resp, err := HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: node.IP, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + AgentVersion: node.AgentVersion, + NginxVersion: node.NginxVersion, + OpenrestyStatus: OpenrestyStatusUnhealthy, + OpenrestyMessage: "port 80 already allocated", }) if err != nil { t.Fatalf("expected heartbeat to succeed: %v", err) @@ -100,6 +103,15 @@ func TestHeartbeatNodeReturnsPreviewUpdateSettings(t *testing.T) { if resp.AgentSettings.UpdateTag != "v0.5.0-rc.1" { t.Fatalf("unexpected update tag: %s", resp.AgentSettings.UpdateTag) } + if !resp.AgentSettings.RestartOpenrestyNow { + t.Fatal("expected restart_openresty_now to be true") + } + if resp.Node.OpenrestyStatus != OpenrestyStatusUnhealthy { + t.Fatalf("expected unhealthy openresty status, got %s", resp.Node.OpenrestyStatus) + } + if resp.Node.OpenrestyMessage != "port 80 already allocated" { + t.Fatalf("unexpected openresty message: %s", resp.Node.OpenrestyMessage) + } storedNode, err := model.GetNodeByID(node.ID) if err != nil { @@ -114,4 +126,24 @@ func TestHeartbeatNodeReturnsPreviewUpdateSettings(t *testing.T) { if storedNode.UpdateTag != "" { t.Fatalf("expected update tag to be cleared, got %s", storedNode.UpdateTag) } + if storedNode.RestartOpenrestyRequested { + t.Fatal("expected restart_openresty_requested to be reset after heartbeat") + } +} + +func TestRequestNodeOpenrestyRestart(t *testing.T) { + setupServiceTestDB(t) + + node, err := CreateNode(NodeInput{Name: "restart-edge-1"}) + if err != nil { + t.Fatalf("failed to create node: %v", err) + } + + updated, err := RequestNodeOpenrestyRestart(node.ID) + if err != nil { + t.Fatalf("expected openresty restart request to succeed: %v", err) + } + if !updated.RestartOpenrestyRequested { + t.Fatal("expected restart_openresty_requested to be true") + } } diff --git a/atsf_server/web/features/nodes/api/nodes.ts b/atsf_server/web/features/nodes/api/nodes.ts index 14731328..257b0cb3 100644 --- a/atsf_server/web/features/nodes/api/nodes.ts +++ b/atsf_server/web/features/nodes/api/nodes.ts @@ -61,3 +61,9 @@ export function requestNodeAgentUpdate( body: JSON.stringify(payload ?? {}), }); } + +export function requestNodeOpenrestyRestart(id: number) { + return apiRequest(`/nodes/${id}/openresty-restart`, { + method: 'POST', + }); +} diff --git a/atsf_server/web/features/nodes/components/node-detail-page.tsx b/atsf_server/web/features/nodes/components/node-detail-page.tsx index b939f0d4..faf667a5 100644 --- a/atsf_server/web/features/nodes/components/node-detail-page.tsx +++ b/atsf_server/web/features/nodes/components/node-detail-page.tsx @@ -21,6 +21,7 @@ import { deleteNode, getNodeAgentRelease, getNodes, + requestNodeOpenrestyRestart, requestNodeAgentUpdate, updateNode, } from '@/features/nodes/api/nodes'; @@ -45,6 +46,8 @@ import { getApplyVariant, getNodeStatusLabel, getNodeStatusVariant, + getOpenrestyStatusLabel, + getOpenrestyStatusVariant, getServerUrl, getUpdateMode, isMeaningfulTime, @@ -202,6 +205,20 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { }, }); + const restartOpenrestyMutation = useMutation({ + mutationFn: () => requestNodeOpenrestyRestart(Number(nodeId)), + onSuccess: async (updatedNode) => { + setFeedback({ + tone: 'success', + message: `已向节点 ${updatedNode.name} 下发 OpenResty 重启指令。`, + }); + await queryClient.invalidateQueries({ queryKey: nodesQueryKey }); + }, + onError: (error) => { + setFeedback({ tone: 'danger', message: getErrorMessage(error) }); + }, + }); + const deleteMutation = useMutation({ mutationFn: () => deleteNode(Number(nodeId)), onSuccess: async () => { @@ -236,6 +253,23 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { deleteMutation.mutate(); }; + const handleRestartOpenresty = () => { + if (!node) { + return; + } + + if ( + !window.confirm( + `确认向节点“${node.name}”下发 OpenResty 重启指令吗?该指令会在下一次心跳后执行。`, + ) + ) { + return; + } + + setFeedback(null); + restartOpenrestyMutation.mutate(); + }; + const handleCopy = async (value: string, message: string) => { try { await copyToClipboard(value); @@ -337,6 +371,17 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { > {node.update_requested ? '查看 Agent 更新' : 'Agent 更新'} + + {restartOpenrestyMutation.isPending + ? '下发重启中...' + : node.restart_openresty_requested + ? '等待 OpenResty 重启' + : '重启 OpenResty'} + + +
+ +

+ {node.restart_openresty_requested + ? '已等待节点在下一次心跳后执行 OpenResty 重启。' + : node.openresty_message || '当前未上报额外错误。'} +

+
+
+
+
+

+ OpenResty 状态消息 +

+

+ {node.openresty_message || '无'} +

+

创建时间 diff --git a/atsf_server/web/features/nodes/components/nodes-page.tsx b/atsf_server/web/features/nodes/components/nodes-page.tsx index 2c467b28..fa66e7f6 100644 --- a/atsf_server/web/features/nodes/components/nodes-page.tsx +++ b/atsf_server/web/features/nodes/components/nodes-page.tsx @@ -36,6 +36,8 @@ import { getApplyVariant, getNodeStatusLabel, getNodeStatusVariant, + getOpenrestyStatusLabel, + getOpenrestyStatusVariant, isMeaningfulTime, } from '@/features/nodes/utils'; @@ -236,7 +238,8 @@ export function NodesPage() { 节点 状态 - Agent / Nginx + Agent / OpenResty + 运行健康 当前版本 最近应用 最近心跳 @@ -266,6 +269,21 @@ export function NodesPage() { {node.agent_version || 'unknown'} /{' '} {node.nginx_version || 'unknown'} + +

+ +

+ {node.openresty_message || '无额外错误'} +

+
+ {node.current_version || '未应用'} diff --git a/atsf_server/web/features/nodes/types.ts b/atsf_server/web/features/nodes/types.ts index 084b0260..1ba03b9a 100644 --- a/atsf_server/web/features/nodes/types.ts +++ b/atsf_server/web/features/nodes/types.ts @@ -10,8 +10,11 @@ export interface NodeItem { update_requested: boolean; update_channel: ReleaseChannel; update_tag: string; + restart_openresty_requested: boolean; agent_version: string; nginx_version: string; + openresty_status: 'healthy' | 'unhealthy' | 'unknown'; + openresty_message: string; status: 'online' | 'offline' | 'pending'; current_version: string; last_seen_at: string; diff --git a/atsf_server/web/features/nodes/utils.ts b/atsf_server/web/features/nodes/utils.ts index d177928e..349872d4 100644 --- a/atsf_server/web/features/nodes/utils.ts +++ b/atsf_server/web/features/nodes/utils.ts @@ -68,6 +68,30 @@ export function getUpdateMode(node: NodeItem) { return { label: '手动更新', variant: 'info' as const }; } +export function getOpenrestyStatusVariant(status: NodeItem['openresty_status']) { + if (status === 'healthy') { + return 'success'; + } + + if (status === 'unhealthy') { + return 'danger'; + } + + return 'warning'; +} + +export function getOpenrestyStatusLabel(status: NodeItem['openresty_status']) { + if (status === 'healthy') { + return '健康'; + } + + if (status === 'unhealthy') { + return '异常'; + } + + return '未知'; +} + function parseVersionParts(version: string) { const normalized = version.trim().replace(/^v/i, ''); if (!normalized || normalized.toLowerCase() === 'unknown') { diff --git a/docs/design.md b/docs/design.md index 6f43de9f..70ded5e0 100644 --- a/docs/design.md +++ b/docs/design.md @@ -21,7 +21,9 @@ ATSFlare 当前定位为内部自用的反向代理控制面,不面向外部 * 反代规则管理 * 配置预览、发布、激活与回滚 * Agent 注册、心跳、同步、应用结果上报 +* Agent 上报 OpenResty 运行健康状态与错误摘要 * OpenResty 配置写入、校验、reload 与失败回滚 +* Server 向 Agent 下发受限运行指令(当前仅支持 OpenResty 重启) * HTTPS/TLS 路由支持 * 证书托管与域名管理 * 节点管理、节点专属 `agent_token`、全局 `discovery_token` @@ -220,6 +222,7 @@ Agent 接口当前覆盖: * 心跳 * 获取激活版本 * 上报应用结果 +* 通过心跳回传 OpenResty 健康状态并接收受限运行指令 统一约束: diff --git a/docs/development-guidelines.md b/docs/development-guidelines.md index 0472274e..353243c9 100644 --- a/docs/development-guidelines.md +++ b/docs/development-guidelines.md @@ -171,6 +171,7 @@ Agent: 禁止: * 将本地 OpenResty 操作暴露为远程执行接口 +* 用通用 shell/命令执行方式替代受限节点操作接口 * 在日志中打印完整 Token --- @@ -208,8 +209,10 @@ Agent 必须满足: * 先执行 `openresty -t` * 成功后执行 `openresty -s reload` * 失败时自动回滚并上报最终结果 +* 周期性向 Server 回传 OpenResty 当前健康状态与最近运行错误摘要 * 支持自动注册与 Token 置换 * 支持接收 Server 下发运行参数 +* 支持接收 Server 下发的受限运行指令,当前仅允许 OpenResty 重启 * 支持自我更新,但失败不影响心跳与同步 ---