[优化] 接口优化, 心跳请求返回规则摘要

This commit is contained in:
ryan
2026-03-13 11:22:37 +08:00
parent b33923d5f7
commit 3fb4cec99c
23 changed files with 300 additions and 177 deletions
+1 -2
View File
@@ -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
}
+1 -1
View File
@@ -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)
+37 -33
View File
@@ -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
}
+16 -14
View File
@@ -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) {
+1 -7
View File
@@ -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)
@@ -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"])
}
+2 -2
View File
@@ -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)
}
+5 -2
View File
@@ -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) {
+15 -5
View File
@@ -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"`
+70 -13
View File
@@ -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,
+60 -6
View File
@@ -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")
}
}