From 2960a4ddebfe8be4f12e081c3ba5bf4fb71ff3b6 Mon Sep 17 00:00:00 2001 From: Ryan <63696351+Rain-kl@users.noreply.github.com> Date: Wed, 7 Oct 2026 10:27:00 +0800 Subject: [PATCH] fix(agent): retry failed node config applies (#40) --- docs/changelog/index.md | 3 + internal/apps/agent/agent/runner.go | 45 ++++++- internal/apps/agent/agent/sync_retry.go | 158 ++++++++++++++++++++++++ internal/apps/agent/heartbeat/cycle.go | 17 ++- 4 files changed, 217 insertions(+), 6 deletions(-) create mode 100644 internal/apps/agent/agent/sync_retry.go diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 3ebaffb9..129e056f 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -10,6 +10,9 @@ sidebar: false ## [Unreleased] +### 🛠 修复 +- 节点配置应用失败后,Agent 会按指数退避自动强制重试,最长间隔为 5 分钟;成功后停止重试,重启后也会恢复未完成的重试。 + ## [v3.5.6] - 2026-10-01 ### ✨ 新功能 diff --git a/internal/apps/agent/agent/runner.go b/internal/apps/agent/agent/runner.go index b2141483..7c5991eb 100644 --- a/internal/apps/agent/agent/runner.go +++ b/internal/apps/agent/agent/runner.go @@ -62,14 +62,28 @@ type Runner struct { RuntimeManager RuntimeManager WebSocketService WebSocketService + configSyncRetry *configSyncRetry restartOpenrestyNow bool websocketUpgradeEnabled bool } // Run starts the agent's main loop, performing heartbeats and upgrading to WebSocket when available. func (r *Runner) Run(ctx context.Context) error { + if r.SyncService != nil { + r.configSyncRetry = newConfigSyncRetry(r.SyncService, r.recordSyncError) + go r.configSyncRetry.run(ctx) + if snapshot, loadErr := r.StateStore.Load(); loadErr != nil { + slog.Error("load state before scheduling config sync retry failed", "error", loadErr) + } else if strings.TrimSpace(snapshot.BlockedVersion) != "" || strings.TrimSpace(snapshot.BlockedChecksum) != "" { + r.configSyncRetry.schedule(&protocol.ActiveConfigMeta{ + Version: snapshot.BlockedVersion, + Checksum: snapshot.BlockedChecksum, + }) + } + } if r.HeartbeatCycle != nil { r.HeartbeatCycle.RecordSyncError = r.recordSyncError + r.HeartbeatCycle.ConfigSyncResult = r.handleConfigSyncResult } nodeID, err := r.StateStore.EnsureNodeID() if err != nil { @@ -313,10 +327,11 @@ func (r *Runner) handleWebSocketMessage(ctx context.Context, message protocol.WS return false, nil } slog.Debug("agent ws active config received", "version", target.Version, "checksum", target.Checksum, "trigger_sync", true) - if err := r.SyncService.SyncOnce(ctx, &target); err != nil { - r.recordSyncError(err) + err := r.SyncService.SyncOnce(ctx, &target) + if err != nil { slog.Error("agent ws triggered sync failed", "version", target.Version, "error", err) } + r.handleConfigSyncResult(&target, err) return false, nil case protocol.WSMessageTypeForceSyncConfig: var target protocol.ActiveConfigMeta @@ -325,10 +340,11 @@ func (r *Runner) handleWebSocketMessage(ctx context.Context, message protocol.WS return false, nil } slog.Debug("agent ws force sync config received", "version", target.Version, "checksum", target.Checksum, "trigger_sync", true) - if err := r.SyncService.ForceSyncOnce(ctx, &target); err != nil { - r.recordSyncError(err) + err := r.SyncService.ForceSyncOnce(ctx, &target) + if err != nil { slog.Error("agent ws triggered force sync failed", "version", target.Version, "error", err) } + r.handleConfigSyncResult(&target, err) return false, nil case protocol.WSMessageTypeWAFIPGroups: var groups []protocol.WAFIPGroup @@ -350,6 +366,27 @@ func (r *Runner) handleWebSocketMessage(ctx context.Context, message protocol.WS } } +func (r *Runner) handleConfigSyncResult(target *protocol.ActiveConfigMeta, err error) { + if err != nil { + r.recordSyncError(err) + if r.configSyncRetry != nil { + r.configSyncRetry.schedule(target) + } + return + } + if r.configSyncRetry == nil || r.StateStore == nil { + return + } + snapshot, loadErr := r.StateStore.Load() + if loadErr != nil { + slog.Error("load state after config sync failed", "error", loadErr) + return + } + if strings.TrimSpace(snapshot.BlockedVersion) == "" && strings.TrimSpace(snapshot.BlockedChecksum) == "" { + r.configSyncRetry.resolve() + } +} + type webSocketBackoff struct { delays []time.Duration index int diff --git a/internal/apps/agent/agent/sync_retry.go b/internal/apps/agent/agent/sync_retry.go new file mode 100644 index 00000000..f29aedfe --- /dev/null +++ b/internal/apps/agent/agent/sync_retry.go @@ -0,0 +1,158 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package agent + +import ( + "context" + "log/slog" + "strings" + "sync" + "time" + + "github.com/Rain-kl/Wavelet/internal/apps/agent/protocol" +) + +const ( + configSyncRetryInitialDelay = 10 * time.Second + configSyncRetryMaxDelay = 5 * time.Minute +) + +type configSyncRetryEvent struct { + target *protocol.ActiveConfigMeta + resolve bool +} + +type configSyncRetry struct { + syncService SyncService + recordError func(error) + mu sync.Mutex + pending *configSyncRetryEvent + wake chan struct{} +} + +func newConfigSyncRetry(syncService SyncService, recordError func(error)) *configSyncRetry { + return &configSyncRetry{ + syncService: syncService, + recordError: recordError, + wake: make(chan struct{}, 1), + } +} + +func (r *configSyncRetry) schedule(target *protocol.ActiveConfigMeta) { + r.notify(configSyncRetryEvent{target: cloneActiveConfigMeta(target)}) +} + +func (r *configSyncRetry) resolve() { + r.notify(configSyncRetryEvent{resolve: true}) +} + +func (r *configSyncRetry) notify(event configSyncRetryEvent) { + r.mu.Lock() + r.pending = &event + r.mu.Unlock() + select { + case r.wake <- struct{}{}: + default: + } +} + +func (r *configSyncRetry) takePending() (configSyncRetryEvent, bool) { + r.mu.Lock() + defer r.mu.Unlock() + if r.pending == nil { + return configSyncRetryEvent{}, false + } + event := *r.pending + r.pending = nil + return event, true +} + +func (r *configSyncRetry) run(ctx context.Context) { + var target *protocol.ActiveConfigMeta + var timer *time.Timer + var timerC <-chan time.Time + attempt := 0 + pending := false + defer func() { + if timer != nil { + timer.Stop() + } + }() + + for { + select { + case <-ctx.Done(): + return + case <-r.wake: + event, ok := r.takePending() + if !ok { + continue + } + if event.resolve { + target = nil + attempt = 0 + pending = false + if timer != nil { + timer.Stop() + } + timerC = nil + continue + } + if !pending || !sameConfigTarget(target, event.target) { + target = event.target + attempt = 0 + pending = true + if timer != nil { + timer.Stop() + } + timer = time.NewTimer(configSyncRetryDelay(attempt)) + timerC = timer.C + } + case <-timerC: + if err := r.syncService.ForceSyncOnce(ctx, target); err != nil { + attempt++ + delay := configSyncRetryDelay(attempt) + if r.recordError != nil { + r.recordError(err) + } + slog.Warn("agent forced config sync retry failed", "attempt", attempt, "retry_after", delay, "error", err) + timer.Reset(delay) + timerC = timer.C + continue + } + slog.Info("agent forced config sync retry succeeded", "attempt", attempt+1) + target = nil + attempt = 0 + pending = false + timerC = nil + } + } +} + +func configSyncRetryDelay(attempt int) time.Duration { + delay := configSyncRetryInitialDelay + for i := 0; i < attempt && delay < configSyncRetryMaxDelay; i++ { + delay *= 2 + if delay > configSyncRetryMaxDelay { + return configSyncRetryMaxDelay + } + } + return delay +} + +func sameConfigTarget(left, right *protocol.ActiveConfigMeta) bool { + if left == nil || right == nil { + return left == nil && right == nil + } + return strings.TrimSpace(left.Version) == strings.TrimSpace(right.Version) && + strings.TrimSpace(left.Checksum) == strings.TrimSpace(right.Checksum) +} + +func cloneActiveConfigMeta(target *protocol.ActiveConfigMeta) *protocol.ActiveConfigMeta { + if target == nil { + return nil + } + cloned := *target + return &cloned +} diff --git a/internal/apps/agent/heartbeat/cycle.go b/internal/apps/agent/heartbeat/cycle.go index 9f614eab..e9833dcb 100644 --- a/internal/apps/agent/heartbeat/cycle.go +++ b/internal/apps/agent/heartbeat/cycle.go @@ -43,6 +43,7 @@ type Cycle struct { Sync SyncService Updater *updater.Service RecordSyncError func(err error) + ConfigSyncResult func(target *protocol.ActiveConfigMeta, err error) } // Perform executes one complete heartbeat cycle: sends the heartbeat, syncs config, and applies settings. @@ -69,14 +70,16 @@ func (c *Cycle) Perform(ctx context.Context, nodeID string, startup bool, settin c.ApplyWAFIPGroups(ctx, heartbeatResult.WAFIPGroups) if startup { if err = c.Sync.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil { - c.recordSyncError(err) slog.Error("agent startup sync failed", "error", err) } else { slog.Debug("agent startup sync completed") } + c.reportConfigSyncResult(heartbeatResult.ActiveConfig, err) } else if err = c.Sync.SyncOnce(ctx, heartbeatResult.ActiveConfig); err != nil { - c.recordSyncError(err) slog.Error("agent sync failed", "error", err) + c.reportConfigSyncResult(heartbeatResult.ActiveConfig, err) + } else { + c.reportConfigSyncResult(heartbeatResult.ActiveConfig, nil) } if settings != nil { settings.RestartOpenrestyIfNeeded(ctx) @@ -85,6 +88,16 @@ func (c *Cycle) Perform(ctx context.Context, nodeID string, startup bool, settin return changed, nil } +func (c *Cycle) reportConfigSyncResult(target *protocol.ActiveConfigMeta, err error) { + if c.ConfigSyncResult != nil { + c.ConfigSyncResult(target, err) + return + } + if err != nil { + c.recordSyncError(err) + } +} + // NodePayload builds and returns the full NodePayload to be sent in a heartbeat request. func (c *Cycle) NodePayload(ctx context.Context, nodeID string) protocol.NodePayload { snapshot, _ := c.StateStore.Load()