diff --git a/atsf_agent/agent.json b/atsf_agent/agent.json index d77fa632..0503990f 100644 --- a/atsf_agent/agent.json +++ b/atsf_agent/agent.json @@ -4,7 +4,6 @@ "data_dir": "./data", "openresty_container_name": "atsflare-openresty", "openresty_docker_image": "openresty/openresty:alpine", - "heartbeat_interval": 30000, - "sync_interval": 30000, + "heartbeat_interval": 10000, "request_timeout": 10000 } diff --git a/atsf_agent/cmd/agent/main.go b/atsf_agent/cmd/agent/main.go index d7c4e761..ab0cdfad 100644 --- a/atsf_agent/cmd/agent/main.go +++ b/atsf_agent/cmd/agent/main.go @@ -38,7 +38,7 @@ func main() { NginxCertDir: cfg.OpenrestyCertDir, }, ) - log.Printf("agent config loaded: server=%s node=%s ip=%s heartbeat_interval=%s sync_interval=%s route_config=%s cert_dir=%s", cfg.ServerURL, cfg.NodeName, cfg.NodeIP, cfg.HeartbeatInterval, cfg.SyncInterval, cfg.RouteConfigPath, cfg.CertDir) + log.Printf("agent config loaded: server=%s node=%s ip=%s heartbeat_interval=%s route_config=%s cert_dir=%s", cfg.ServerURL, cfg.NodeName, cfg.NodeIP, cfg.HeartbeatInterval, cfg.RouteConfigPath, cfg.CertDir) client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration()) stateStore := state.NewStore(cfg.StatePath) diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index 07734a8a..bcd729cf 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -14,13 +14,13 @@ import ( type HeartbeatService interface { Register(ctx context.Context, payload protocol.NodePayload) (*protocol.RegisterNodeResponse, error) - Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.AgentSettings, error) + Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.HeartbeatResult, error) SetToken(token string) } type SyncService interface { - SyncOnStartup(ctx context.Context) error - SyncOnce(ctx context.Context) error + SyncOnStartup(ctx context.Context, target *protocol.ActiveConfigMeta) error + SyncOnce(ctx context.Context, target *protocol.ActiveConfigMeta) error } type Updater interface { @@ -61,20 +61,24 @@ func (r *Runner) Run(ctx context.Context) error { } log.Printf("agent runner started: node_id=%s node=%s ip=%s", nodeID, r.Config.NodeName, r.Config.NodeIP) if r.hasAgentToken() { - if err = r.SyncService.SyncOnStartup(ctx); err != nil { - r.recordSyncError(err) - log.Printf("agent startup sync failed: %v", err) - } else { - log.Printf("agent startup sync completed") - } r.refreshOpenrestyHealth(ctx) - settings, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) + heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) if hbErr != nil { log.Printf("agent startup heartbeat failed: %v", hbErr) } else { + if heartbeatResult == nil { + heartbeatResult = &protocol.HeartbeatResult{} + } log.Printf("agent startup heartbeat succeeded: node_id=%s", nodeID) - r.applySettings(settings) + r.applySettings(heartbeatResult.AgentSettings) + if err = r.SyncService.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil { + r.recordSyncError(err) + log.Printf("agent startup sync failed: %v", err) + } else { + log.Printf("agent startup sync completed") + } r.tryRestartOpenresty(ctx) + r.tryAutoUpdate(ctx) } } else if err = r.tryRegister(ctx, &nodeID); err != nil { log.Printf("agent initial discovery register failed: %v", err) @@ -82,8 +86,6 @@ func (r *Runner) Run(ctx context.Context) error { heartbeatTicker := time.NewTicker(r.Config.HeartbeatInterval.Duration()) defer heartbeatTicker.Stop() - syncTicker := time.NewTicker(r.Config.SyncInterval.Duration()) - defer syncTicker.Stop() for { select { @@ -98,25 +100,23 @@ func (r *Runner) Run(ctx context.Context) error { continue } r.refreshOpenrestyHealth(ctx) - settings, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) + heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) if hbErr != nil { log.Printf("agent heartbeat failed: %v", hbErr) } else { - if changed := r.applySettings(settings); changed { + if heartbeatResult == nil { + heartbeatResult = &protocol.HeartbeatResult{} + } + if changed := r.applySettings(heartbeatResult.AgentSettings); changed { heartbeatTicker.Reset(r.Config.HeartbeatInterval.Duration()) - syncTicker.Reset(r.Config.SyncInterval.Duration()) + } + if err = r.SyncService.SyncOnce(ctx, heartbeatResult.ActiveConfig); err != nil { + r.recordSyncError(err) + log.Printf("agent sync failed: %v", err) } r.tryRestartOpenresty(ctx) r.tryAutoUpdate(ctx) } - case <-syncTicker.C: - if !r.hasAgentToken() { - continue - } - if err = r.SyncService.SyncOnce(ctx); err != nil { - r.recordSyncError(err) - log.Printf("agent sync failed: %v", err) - } } } } @@ -138,14 +138,6 @@ func (r *Runner) applySettings(settings *protocol.AgentSettings) bool { changed = true } } - if settings.SyncInterval > 0 { - newInterval := config.MillisecondDuration(time.Duration(settings.SyncInterval) * time.Millisecond) - if newInterval != r.Config.SyncInterval { - log.Printf("agent sync interval updated: %s -> %s", r.Config.SyncInterval, newInterval) - r.Config.SyncInterval = newInterval - changed = true - } - } r.autoUpdate = settings.AutoUpdate r.updateNow = settings.UpdateNow r.updateRepo = strings.TrimSpace(settings.UpdateRepo) @@ -226,12 +218,24 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error { r.HeartbeatService.SetToken(response.AgentToken) *nodeID = response.NodeID log.Printf("agent discovery registration succeeded: node_id=%s", response.NodeID) - if err = r.SyncService.SyncOnStartup(ctx); err != nil { + r.refreshOpenrestyHealth(ctx) + heartbeatResult, heartbeatErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(*nodeID)) + if heartbeatErr != nil { + log.Printf("agent post-register heartbeat failed: %v", heartbeatErr) + return nil + } + if heartbeatResult == nil { + heartbeatResult = &protocol.HeartbeatResult{} + } + r.applySettings(heartbeatResult.AgentSettings) + if err = r.SyncService.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil { r.recordSyncError(err) log.Printf("agent post-register startup sync failed: %v", err) } else { log.Printf("agent post-register startup sync completed") } + r.tryRestartOpenresty(ctx) + r.tryAutoUpdate(ctx) return nil } diff --git a/atsf_agent/internal/agent/runner_test.go b/atsf_agent/internal/agent/runner_test.go index 4dad75f2..c4e35c67 100644 --- a/atsf_agent/internal/agent/runner_test.go +++ b/atsf_agent/internal/agent/runner_test.go @@ -21,7 +21,7 @@ type fakeHeartbeatService struct { registerErr error registerResp *protocol.RegisterNodeResponse heartbeatErrs []error - heartbeatSettings []*protocol.AgentSettings + heartbeatResults []*protocol.HeartbeatResult heartbeatPayloads []protocol.NodePayload onHeartbeat func(int) lastToken string @@ -34,7 +34,7 @@ func (f *fakeHeartbeatService) Register(ctx context.Context, payload protocol.No return f.registerResp, f.registerErr } -func (f *fakeHeartbeatService) Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.AgentSettings, error) { +func (f *fakeHeartbeatService) Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.HeartbeatResult, error) { f.mu.Lock() f.heartbeatCalls++ callIndex := f.heartbeatCalls @@ -43,16 +43,16 @@ func (f *fakeHeartbeatService) Heartbeat(ctx context.Context, payload protocol.N if len(f.heartbeatErrs) >= callIndex { err = f.heartbeatErrs[callIndex-1] } - var settings *protocol.AgentSettings - if len(f.heartbeatSettings) >= callIndex { - settings = f.heartbeatSettings[callIndex-1] + var result *protocol.HeartbeatResult + if len(f.heartbeatResults) >= callIndex { + result = f.heartbeatResults[callIndex-1] } onHeartbeat := f.onHeartbeat f.mu.Unlock() if onHeartbeat != nil { onHeartbeat(callIndex) } - return settings, err + return result, err } func (f *fakeHeartbeatService) SetToken(token string) { @@ -94,14 +94,14 @@ func (f *fakeRuntimeManager) Restart(ctx context.Context) error { return f.restartErr } -func (f *fakeSyncService) SyncOnStartup(ctx context.Context) error { +func (f *fakeSyncService) SyncOnStartup(ctx context.Context, target *protocol.ActiveConfigMeta) error { f.mu.Lock() defer f.mu.Unlock() f.startupCalls++ return f.startupErr } -func (f *fakeSyncService) SyncOnce(ctx context.Context) error { +func (f *fakeSyncService) SyncOnce(ctx context.Context, target *protocol.ActiveConfigMeta) error { f.mu.Lock() f.syncOnceCalls++ callIndex := f.syncOnceCalls @@ -119,6 +119,7 @@ func TestRunnerKeepsHeartbeatWhenStartupSyncFails(t *testing.T) { stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) heartbeatService := &fakeHeartbeatService{ + heartbeatResults: []*protocol.HeartbeatResult{{}}, onHeartbeat: func(callCount int) { if callCount >= 2 { cancel() @@ -136,7 +137,6 @@ func TestRunnerKeepsHeartbeatWhenStartupSyncFails(t *testing.T) { AgentVersion: config.AgentVersion, NginxVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), - SyncInterval: config.MillisecondDuration(20 * time.Millisecond), }, StateStore: stateStore, HeartbeatService: heartbeatService, @@ -170,6 +170,9 @@ func TestRunnerDoesNotExitOnHeartbeatOrSyncError(t *testing.T) { heartbeatService := &fakeHeartbeatService{ registerErr: errors.New("register timeout"), heartbeatErrs: []error{errors.New("heartbeat timeout")}, + heartbeatResults: []*protocol.HeartbeatResult{ + {}, + }, } syncService := &fakeSyncService{ syncOnceErr: errors.New("openresty reload failed"), @@ -187,7 +190,6 @@ func TestRunnerDoesNotExitOnHeartbeatOrSyncError(t *testing.T) { AgentVersion: config.AgentVersion, NginxVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), - SyncInterval: config.MillisecondDuration(10 * time.Millisecond), }, StateStore: stateStore, HeartbeatService: heartbeatService, @@ -225,7 +227,9 @@ func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) { t.Fatalf("failed to seed state: %v", err) } heartbeatService := &fakeHeartbeatService{ - heartbeatSettings: []*protocol.AgentSettings{{RestartOpenrestyNow: true}}, + heartbeatResults: []*protocol.HeartbeatResult{{ + AgentSettings: &protocol.AgentSettings{RestartOpenrestyNow: true}, + }}, onHeartbeat: func(callCount int) { if callCount >= 1 { cancel() @@ -244,7 +248,6 @@ func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) { AgentVersion: config.AgentVersion, NginxVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), - SyncInterval: config.MillisecondDuration(100 * time.Millisecond), }, StateStore: stateStore, HeartbeatService: heartbeatService, @@ -289,6 +292,7 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { AgentToken: "agent-token-issued", Name: "edge-01", }, + heartbeatResults: []*protocol.HeartbeatResult{{}}, onHeartbeat: func(callCount int) { if callCount >= 1 { cancel() @@ -313,7 +317,6 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { AgentVersion: config.AgentVersion, NginxVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), - SyncInterval: config.MillisecondDuration(20 * time.Millisecond), }, StateStore: stateStore, HeartbeatService: heartbeatService, @@ -323,7 +326,6 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { runner.Config.AgentVersion = config.AgentVersion runner.Config.NginxVersion = "1.27.1.2" runner.Config.HeartbeatInterval = config.MillisecondDuration(10 * time.Millisecond) - runner.Config.SyncInterval = config.MillisecondDuration(20 * time.Millisecond) err = runner.Run(ctx) if !errors.Is(err, context.Canceled) { diff --git a/atsf_agent/internal/config/config.go b/atsf_agent/internal/config/config.go index b5b00533..78763bae 100644 --- a/atsf_agent/internal/config/config.go +++ b/atsf_agent/internal/config/config.go @@ -38,7 +38,6 @@ type Config struct { OpenrestyCertDir string `json:"openresty_cert_dir"` StatePath string `json:"state_path"` HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` - SyncInterval MillisecondDuration `json:"sync_interval"` RequestTimeout MillisecondDuration `json:"request_timeout"` configPath string } @@ -60,7 +59,6 @@ type configFile struct { OpenrestyCertDir string `json:"openresty_cert_dir"` StatePath string `json:"state_path"` HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` - SyncInterval MillisecondDuration `json:"sync_interval"` RequestTimeout MillisecondDuration `json:"request_timeout"` } @@ -90,7 +88,6 @@ func Load(path string) (*Config, error) { OpenrestyCertDir: file.OpenrestyCertDir, StatePath: file.StatePath, HeartbeatInterval: file.HeartbeatInterval, - SyncInterval: file.SyncInterval, RequestTimeout: file.RequestTimeout, } cfg.configPath = path @@ -148,10 +145,7 @@ func applyDefaults(cfg *Config, baseDir string) { } } if cfg.HeartbeatInterval <= 0 { - cfg.HeartbeatInterval = MillisecondDuration(30 * time.Second) - } - if cfg.SyncInterval <= 0 { - cfg.SyncInterval = MillisecondDuration(30 * time.Second) + cfg.HeartbeatInterval = MillisecondDuration(10 * time.Second) } if cfg.RequestTimeout <= 0 { cfg.RequestTimeout = MillisecondDuration(10 * time.Second) diff --git a/atsf_agent/internal/config/config_test.go b/atsf_agent/internal/config/config_test.go index 34c9c32a..9dd37704 100644 --- a/atsf_agent/internal/config/config_test.go +++ b/atsf_agent/internal/config/config_test.go @@ -142,7 +142,6 @@ func TestLoadUsesMillisecondsForIntervals(t *testing.T) { "node_name": "edge-01", "node_ip": "10.0.0.8", "heartbeat_interval": 30000, - "sync_interval": 45000, "request_timeout": 1500, } data, err := json.Marshal(payload) @@ -161,9 +160,6 @@ func TestLoadUsesMillisecondsForIntervals(t *testing.T) { if cfg.HeartbeatInterval.Duration() != 30*time.Second { t.Fatalf("unexpected heartbeat interval: %s", cfg.HeartbeatInterval) } - if cfg.SyncInterval.Duration() != 45*time.Second { - t.Fatalf("unexpected sync interval: %s", cfg.SyncInterval) - } if cfg.RequestTimeout.Duration() != 1500*time.Millisecond { t.Fatalf("unexpected request timeout: %s", cfg.RequestTimeout) } @@ -182,7 +178,6 @@ func TestSavePersistsMillisecondsAndOmitsRuntimeVersions(t *testing.T) { } cfg.NginxVersion = "1.27.1.2" cfg.HeartbeatInterval = MillisecondDuration(5 * time.Second) - cfg.SyncInterval = MillisecondDuration(6 * time.Second) cfg.RequestTimeout = MillisecondDuration(7 * time.Second) if err = cfg.Save(); err != nil { @@ -206,9 +201,6 @@ func TestSavePersistsMillisecondsAndOmitsRuntimeVersions(t *testing.T) { if decoded["heartbeat_interval"] != float64(5000) { t.Fatalf("unexpected heartbeat interval: %#v", decoded["heartbeat_interval"]) } - if decoded["sync_interval"] != float64(6000) { - t.Fatalf("unexpected sync interval: %#v", decoded["sync_interval"]) - } if decoded["request_timeout"] != float64(7000) { t.Fatalf("unexpected request timeout: %#v", decoded["request_timeout"]) } diff --git a/atsf_agent/internal/heartbeat/service.go b/atsf_agent/internal/heartbeat/service.go index a9e733d1..9ad8cc8d 100644 --- a/atsf_agent/internal/heartbeat/service.go +++ b/atsf_agent/internal/heartbeat/service.go @@ -8,7 +8,7 @@ import ( type Client interface { RegisterNode(ctx context.Context, payload protocol.NodePayload) (*protocol.RegisterNodeResponse, error) - Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.AgentSettings, error) + Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.HeartbeatResult, error) SetToken(token string) } @@ -24,7 +24,7 @@ func (s *Service) Register(ctx context.Context, payload protocol.NodePayload) (* return s.client.RegisterNode(ctx, payload) } -func (s *Service) Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.AgentSettings, error) { +func (s *Service) Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.HeartbeatResult, error) { return s.client.Heartbeat(ctx, payload) } diff --git a/atsf_agent/internal/httpclient/client.go b/atsf_agent/internal/httpclient/client.go index 265f65d4..dac883a0 100644 --- a/atsf_agent/internal/httpclient/client.go +++ b/atsf_agent/internal/httpclient/client.go @@ -42,7 +42,7 @@ func (c *Client) RegisterNode(ctx context.Context, payload protocol.NodePayload) return &resp.Data, nil } -func (c *Client) Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.AgentSettings, error) { +func (c *Client) Heartbeat(ctx context.Context, payload protocol.NodePayload) (*protocol.HeartbeatResult, error) { resp := protocol.HeartbeatAPIResponse{} if err := c.postJSON(ctx, "/api/agent/nodes/heartbeat", payload, &resp); err != nil { return nil, err @@ -50,7 +50,10 @@ func (c *Client) Heartbeat(ctx context.Context, payload protocol.NodePayload) (* if !resp.Success { return nil, errors.New(resp.Message) } - return resp.AgentSettings, nil + return &protocol.HeartbeatResult{ + AgentSettings: resp.AgentSettings, + ActiveConfig: resp.ActiveConfig, + }, nil } func (c *Client) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error) { diff --git a/atsf_agent/internal/protocol/agent_api.go b/atsf_agent/internal/protocol/agent_api.go index 24d8d86f..931df306 100644 --- a/atsf_agent/internal/protocol/agent_api.go +++ b/atsf_agent/internal/protocol/agent_api.go @@ -7,15 +7,20 @@ type APIResponse[T any] struct { } type HeartbeatAPIResponse struct { - Success bool `json:"success"` - Message string `json:"message"` - Data any `json:"data"` - AgentSettings *AgentSettings `json:"agent_settings,omitempty"` + Success bool `json:"success"` + Message string `json:"message"` + Data any `json:"data"` + AgentSettings *AgentSettings `json:"agent_settings,omitempty"` + ActiveConfig *ActiveConfigMeta `json:"active_config,omitempty"` +} + +type HeartbeatResult struct { + AgentSettings *AgentSettings + ActiveConfig *ActiveConfigMeta } 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"` @@ -69,6 +74,11 @@ type ActiveConfigResponse struct { CreatedAt string `json:"created_at"` } +type ActiveConfigMeta struct { + Version string `json:"version"` + Checksum string `json:"checksum"` +} + type SupportFile struct { Path string `json:"path"` Content string `json:"content"` diff --git a/atsf_agent/internal/sync/service.go b/atsf_agent/internal/sync/service.go index 08418025..dff75b28 100644 --- a/atsf_agent/internal/sync/service.go +++ b/atsf_agent/internal/sync/service.go @@ -5,6 +5,7 @@ import ( "crypto/sha256" "encoding/hex" "log" + "strings" "atsflare-agent/internal/protocol" "atsflare-agent/internal/state" @@ -40,15 +41,15 @@ func New(client ConfigClient, nginxManager NginxManager, stateStore *state.Store } } -func (s *Service) SyncOnce(ctx context.Context) error { - return s.sync(ctx, false) +func (s *Service) SyncOnce(ctx context.Context, target *protocol.ActiveConfigMeta) error { + return s.sync(ctx, false, target) } -func (s *Service) SyncOnStartup(ctx context.Context) error { - return s.sync(ctx, true) +func (s *Service) SyncOnStartup(ctx context.Context, target *protocol.ActiveConfigMeta) error { + return s.sync(ctx, true, target) } -func (s *Service) sync(ctx context.Context, startup bool) error { +func (s *Service) sync(ctx context.Context, startup bool, target *protocol.ActiveConfigMeta) error { mode := "periodic" if startup { mode = "startup" @@ -57,20 +58,73 @@ func (s *Service) sync(ctx context.Context, startup bool) error { if err != nil { return err } + currentChecksum, err := s.nginxManager.CurrentChecksum() + if err != nil { + return err + } + + if target != nil { + target.Version = strings.TrimSpace(target.Version) + target.Checksum = strings.TrimSpace(target.Checksum) + } + + if target == nil || target.Version == "" || target.Checksum == "" { + if !startup { + log.Printf("skipping sync because heartbeat returned no active config summary: mode=%s", mode) + return nil + } + log.Printf("sync startup fallback: active config summary unavailable, fetching active config directly") + config, fetchErr := s.client.GetActiveConfig(ctx) + if fetchErr != nil { + log.Printf("fetch active config failed: mode=%s error=%v", mode, fetchErr) + return fetchErr + } + target = &protocol.ActiveConfigMeta{ + Version: config.Version, + Checksum: config.Checksum, + } + return s.applyIfNeeded(ctx, mode, startup, snapshot, currentChecksum, target, config) + } + + if currentChecksum == target.Checksum { + log.Printf("local openresty config already up to date: mode=%s version=%s", mode, target.Version) + if startup { + log.Printf("ensuring openresty runtime on startup: version=%s", target.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", target.Version) + snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy + snapshot.OpenrestyMessage = "" + } + snapshot.CurrentVersion = target.Version + snapshot.CurrentChecksum = target.Checksum + snapshot.LastError = "" + log.Printf("sync finished without changes: mode=%s version=%s", mode, target.Version) + return s.stateStore.Save(snapshot) + } + if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup { + log.Printf("skipping config fetch because state already records target version/checksum: version=%s checksum=%s", target.Version, target.Checksum) + return nil + } + config, err := s.client.GetActiveConfig(ctx) if err != nil { log.Printf("fetch active config failed: mode=%s error=%v", mode, err) return err } - currentChecksum, err := s.nginxManager.CurrentChecksum() - if err != nil { - return err - } + return s.applyIfNeeded(ctx, mode, startup, snapshot, currentChecksum, target, config) +} + +func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool, snapshot *state.Snapshot, currentChecksum string, target *protocol.ActiveConfigMeta, config *protocol.ActiveConfigResponse) error { if currentChecksum == config.Checksum { log.Printf("local openresty config already up to date: mode=%s version=%s", mode, config.Version) if startup { log.Printf("ensuring openresty runtime on startup: version=%s", config.Version) - if err = s.nginxManager.EnsureRuntime(ctx, true); err != nil { + if err := s.nginxManager.EnsureRuntime(ctx, true); err != nil { snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy snapshot.OpenrestyMessage = err.Error() _ = s.stateStore.Save(snapshot) @@ -86,6 +140,9 @@ func (s *Service) sync(ctx context.Context, startup bool) error { log.Printf("sync finished without changes: mode=%s version=%s", mode, config.Version) return s.stateStore.Save(snapshot) } + if target != nil && (target.Version != config.Version || target.Checksum != config.Checksum) { + log.Printf("active config changed between heartbeat and fetch: heartbeat_version=%s heartbeat_checksum=%s fetched_version=%s fetched_checksum=%s", target.Version, target.Checksum, config.Version, config.Checksum) + } if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum && !startup { log.Printf("skipping apply because state already records target version/checksum: version=%s checksum=%s", config.Version, config.Checksum) return nil @@ -97,7 +154,7 @@ func (s *Service) sync(ctx context.Context, startup bool) error { mainConfigChecksum := checksumString(config.MainConfig) routeConfigChecksum := checksumString(routeConfig) log.Printf("applying new openresty config: mode=%s from_version=%s to_version=%s old_checksum=%s new_checksum=%s", mode, snapshot.CurrentVersion, config.Version, currentChecksum, config.Checksum) - if err = s.nginxManager.Apply(ctx, config.MainConfig, routeConfig, config.SupportFiles); err != nil { + if err := s.nginxManager.Apply(ctx, config.MainConfig, routeConfig, 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 @@ -126,10 +183,10 @@ func (s *Service) sync(ctx context.Context, startup bool) error { snapshot.LastError = "" snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy snapshot.OpenrestyMessage = "" - if err = s.stateStore.Save(snapshot); err != nil { + if err := s.stateStore.Save(snapshot); err != nil { return err } - if err = s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{ + if err := s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{ NodeID: snapshot.NodeID, Version: config.Version, Result: ApplyResultSuccess, diff --git a/atsf_agent/internal/sync/service_test.go b/atsf_agent/internal/sync/service_test.go index 67538f06..c9629f1f 100644 --- a/atsf_agent/internal/sync/service_test.go +++ b/atsf_agent/internal/sync/service_test.go @@ -18,8 +18,9 @@ type fakeExecutor struct { } type fakeClient struct { - config protocol.ActiveConfigResponse - reports []protocol.ApplyLogPayload + config protocol.ActiveConfigResponse + reports []protocol.ApplyLogPayload + fetchCalls int } type fakeManager struct { @@ -54,6 +55,7 @@ func (f *fakeExecutor) Restart(ctx context.Context) error { } func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error) { + f.fetchCalls++ return &f.config, nil } @@ -109,7 +111,10 @@ func TestSyncOnceSuccess(t *testing.T) { Executor: &fakeExecutor{}, }, stateStore) - if err = service.SyncOnce(context.Background()); err != nil { + if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{ + Version: client.config.Version, + Checksum: client.config.Checksum, + }); err != nil { t.Fatalf("SyncOnce failed: %v", err) } @@ -192,7 +197,10 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) { }, }, stateStore) - err = service.SyncOnce(context.Background()) + err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{ + Version: client.config.Version, + Checksum: client.config.Checksum, + }) if err == nil { t.Fatal("expected SyncOnce to fail when nginx test fails") } @@ -255,7 +263,10 @@ func TestSyncOnStartupRecreatesRuntimeWhenChecksumMatches(t *testing.T) { manager := &fakeManager{currentChecksum: "checksum-3"} service := New(client, manager, stateStore) - if err = service.SyncOnStartup(context.Background()); err != nil { + if err = service.SyncOnStartup(context.Background(), &protocol.ActiveConfigMeta{ + Version: client.config.Version, + Checksum: client.config.Checksum, + }); err != nil { t.Fatalf("SyncOnStartup failed: %v", err) } if len(manager.ensureCalls) != 1 || !manager.ensureCalls[0] { @@ -301,7 +312,10 @@ func TestSyncOnStartupRecordsRuntimeFailure(t *testing.T) { ensureErr: context.DeadlineExceeded, } service := New(client, manager, stateStore) - if err = service.SyncOnStartup(context.Background()); err == nil { + if err = service.SyncOnStartup(context.Background(), &protocol.ActiveConfigMeta{ + Version: client.config.Version, + Checksum: client.config.Checksum, + }); err == nil { t.Fatal("expected SyncOnStartup to fail when runtime recreation fails") } snapshot, err := stateStore.Load() @@ -315,3 +329,43 @@ func TestSyncOnStartupRecordsRuntimeFailure(t *testing.T) { t.Fatal("expected runtime error message to be recorded") } } + +func TestSyncOnceSkipsFetchWhenHeartbeatChecksumMatches(t *testing.T) { + client := &fakeClient{ + config: protocol.ActiveConfigResponse{ + Version: "20260309-005", + Checksum: "checksum-5", + MainConfig: "worker_processes auto;", + RouteConfig: "server { listen 84; }", + RenderedConfig: "server { listen 84; }", + 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, + CurrentVersion: client.config.Version, + CurrentChecksum: client.config.Checksum, + }); err != nil { + t.Fatalf("failed to seed state: %v", err) + } + + manager := &fakeManager{currentChecksum: client.config.Checksum} + service := New(client, manager, stateStore) + if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{ + Version: client.config.Version, + Checksum: client.config.Checksum, + }); err != nil { + t.Fatalf("SyncOnce failed: %v", err) + } + if client.fetchCalls != 0 { + t.Fatalf("expected no active config fetch when heartbeat checksum matches, got %d", client.fetchCalls) + } + if len(client.reports) != 0 { + t.Fatal("expected no apply log when no config change is needed") + } +} diff --git a/atsf_server/common/constants.go b/atsf_server/common/constants.go index 05f5bdb8..e2bcecf2 100644 --- a/atsf_server/common/constants.go +++ b/atsf_server/common/constants.go @@ -51,8 +51,7 @@ var AgentDiscoveryToken = "" var NodeOfflineThreshold = 2 * time.Minute // V3 operational settings (hot-reloadable via Option table) -var AgentHeartbeatInterval = 30000 // milliseconds -var AgentSyncInterval = 30000 // milliseconds +var AgentHeartbeatInterval = 10000 // milliseconds var AgentUpdateRepo = "Rain-kl/ATSFlare" // V5 OpenResty performance settings (hot-reloadable via Option table) diff --git a/atsf_server/controller/agent.go b/atsf_server/controller/agent.go index fa235d73..95a24dcf 100644 --- a/atsf_server/controller/agent.go +++ b/atsf_server/controller/agent.go @@ -91,6 +91,7 @@ func AgentHeartbeat(c *gin.Context) { "message": "", "data": node.Node, "agent_settings": node.AgentSettings, + "active_config": node.ActiveConfig, }) } diff --git a/atsf_server/model/option.go b/atsf_server/model/option.go index 93c5b175..5be3dbd6 100644 --- a/atsf_server/model/option.go +++ b/atsf_server/model/option.go @@ -52,7 +52,6 @@ func InitOptionMap() { common.OptionMap["TurnstileSecretKey"] = "" common.OptionMap["AgentDiscoveryToken"] = "" common.OptionMap["AgentHeartbeatInterval"] = strconv.Itoa(common.AgentHeartbeatInterval) - common.OptionMap["AgentSyncInterval"] = strconv.Itoa(common.AgentSyncInterval) common.OptionMap["NodeOfflineThreshold"] = strconv.Itoa(int(common.NodeOfflineThreshold.Milliseconds())) common.OptionMap["AgentUpdateRepo"] = common.AgentUpdateRepo common.OptionMap["OpenRestyWorkerProcesses"] = common.OpenRestyWorkerProcesses @@ -199,10 +198,6 @@ func updateOptionMap(key string, value string) { if v, err := strconv.Atoi(value); err == nil && v > 0 { common.AgentHeartbeatInterval = v } - case "AgentSyncInterval": - if v, err := strconv.Atoi(value); err == nil && v > 0 { - common.AgentSyncInterval = v - } case "NodeOfflineThreshold": if v, err := strconv.Atoi(value); err == nil && v > 0 { common.NodeOfflineThreshold = time.Duration(v) * time.Millisecond diff --git a/atsf_server/router/api_phase2_test.go b/atsf_server/router/api_phase2_test.go index c373c5a3..c3807fb0 100644 --- a/atsf_server/router/api_phase2_test.go +++ b/atsf_server/router/api_phase2_test.go @@ -284,9 +284,10 @@ func TestPhase2AgentLifecycle(t *testing.T) { 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"` + Success bool `json:"success"` + Message string `json:"message"` + AgentSettings service.AgentSettings `json:"agent_settings"` + ActiveConfig *service.ActiveConfigMeta `json:"active_config"` } if err = json.Unmarshal(restartHeartbeatRecorder.Body.Bytes(), &restartHeartbeatBody); err != nil { t.Fatalf("failed to decode heartbeat response: %v", err) @@ -297,6 +298,9 @@ func TestPhase2AgentLifecycle(t *testing.T) { if !restartHeartbeatBody.AgentSettings.RestartOpenrestyNow { t.Fatal("expected heartbeat response to instruct openresty restart") } + if restartHeartbeatBody.ActiveConfig == nil || restartHeartbeatBody.ActiveConfig.Version == "" || restartHeartbeatBody.ActiveConfig.Checksum == "" { + t.Fatal("expected heartbeat response to include active config summary") + } 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 0de02339..773682bf 100644 --- a/atsf_server/service/agent.go +++ b/atsf_server/service/agent.go @@ -57,7 +57,6 @@ 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"` @@ -66,9 +65,15 @@ type AgentSettings struct { RestartOpenrestyNow bool `json:"restart_openresty_now"` } +type ActiveConfigMeta struct { + Version string `json:"version"` + Checksum string `json:"checksum"` +} + type HeartbeatResponse struct { - Node *model.Node `json:"node"` - AgentSettings *AgentSettings `json:"agent_settings"` + Node *model.Node `json:"node"` + AgentSettings *AgentSettings `json:"agent_settings"` + ActiveConfig *ActiveConfigMeta `json:"active_config"` } type NodeView struct { @@ -124,11 +129,14 @@ func HeartbeatNode(node *model.Node, payload AgentNodePayload) (*HeartbeatRespon 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 } + activeConfig, err := GetActiveConfigMetaForAgent() + if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { + return nil, err + } return &HeartbeatResponse{ Node: node, AgentSettings: &AgentSettings{ HeartbeatInterval: common.AgentHeartbeatInterval, - SyncInterval: common.AgentSyncInterval, AutoUpdate: node.AutoUpdateEnabled, UpdateRepo: common.AgentUpdateRepo, UpdateNow: updateNow, @@ -136,6 +144,21 @@ func HeartbeatNode(node *model.Node, payload AgentNodePayload) (*HeartbeatRespon UpdateTag: updateTag, RestartOpenrestyNow: restartOpenrestyNow, }, + ActiveConfig: activeConfig, + }, nil +} + +func GetActiveConfigMetaForAgent() (*ActiveConfigMeta, error) { + version, err := model.GetActiveConfigVersion() + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, err + } + return nil, err + } + return &ActiveConfigMeta{ + Version: version.Version, + Checksum: version.Checksum, }, nil } diff --git a/atsf_server/service/node_update_test.go b/atsf_server/service/node_update_test.go index fb267b73..f6a6d914 100644 --- a/atsf_server/service/node_update_test.go +++ b/atsf_server/service/node_update_test.go @@ -78,6 +78,17 @@ func TestHeartbeatNodeReturnsPreviewUpdateSettings(t *testing.T) { if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) } + if err := model.DB.Create(&model.ConfigVersion{ + Version: "20260313-001", + SnapshotJSON: "{}", + MainConfig: "worker_processes auto;", + RenderedConfig: "server { listen 80; }", + Checksum: "checksum-active-1", + IsActive: true, + CreatedBy: "root", + }).Error; err != nil { + t.Fatalf("failed to seed active config version: %v", err) + } resp, err := HeartbeatNode(node, AgentNodePayload{ NodeID: node.NodeID, @@ -94,6 +105,12 @@ func TestHeartbeatNodeReturnsPreviewUpdateSettings(t *testing.T) { if resp.AgentSettings == nil { t.Fatal("expected agent settings in heartbeat response") } + if resp.ActiveConfig == nil { + t.Fatal("expected active config summary in heartbeat response") + } + if resp.ActiveConfig.Version == "" || resp.ActiveConfig.Checksum == "" { + t.Fatal("expected active config summary to include version and checksum") + } if !resp.AgentSettings.UpdateNow { t.Fatal("expected update_now to be true") } diff --git a/atsf_server/web/features/settings/components/settings-page.tsx b/atsf_server/web/features/settings/components/settings-page.tsx index 55c2fcf4..875d2c96 100644 --- a/atsf_server/web/features/settings/components/settings-page.tsx +++ b/atsf_server/web/features/settings/components/settings-page.tsx @@ -68,8 +68,7 @@ const defaultSystemFields = { }; const defaultOperationFields = { - AgentHeartbeatInterval: '30000', - AgentSyncInterval: '30000', + AgentHeartbeatInterval: '10000', NodeOfflineThreshold: '120000', AgentUpdateRepo: 'Rain-kl/ATSFlare', OpenRestyWorkerProcesses: 'auto', @@ -335,8 +334,7 @@ export function SettingsPage() { }); setOperationFields({ - AgentHeartbeatInterval: optionMap.AgentHeartbeatInterval ?? '30000', - AgentSyncInterval: optionMap.AgentSyncInterval ?? '30000', + AgentHeartbeatInterval: optionMap.AgentHeartbeatInterval ?? '10000', NodeOfflineThreshold: optionMap.NodeOfflineThreshold ?? '120000', AgentUpdateRepo: optionMap.AgentUpdateRepo ?? 'Rain-kl/ATSFlare', OpenRestyWorkerProcesses: optionMap.OpenRestyWorkerProcesses ?? 'auto', @@ -942,10 +940,6 @@ export function SettingsPage() { operationFields.AgentHeartbeatInterval, 10, ); - const sync = Number.parseInt( - operationFields.AgentSyncInterval, - 10, - ); const offline = Number.parseInt( operationFields.NodeOfflineThreshold, 10, @@ -954,9 +948,6 @@ export function SettingsPage() { if (Number.isNaN(heartbeat) || heartbeat < 5000) { throw new Error('心跳间隔不能小于 5000 毫秒。'); } - if (Number.isNaN(sync) || sync < 5000) { - throw new Error('同步间隔不能小于 5000 毫秒。'); - } if (Number.isNaN(offline) || offline < 10000) { throw new Error('离线阈值不能小于 10000 毫秒。'); } @@ -964,7 +955,6 @@ export function SettingsPage() { await saveOptionEntries( [ ['AgentHeartbeatInterval', String(heartbeat)], - ['AgentSyncInterval', String(sync)], ['NodeOfflineThreshold', String(offline)], ], 'Agent 运行参数已保存。', @@ -980,7 +970,7 @@ export function SettingsPage() { } > -