mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-08 00:26:37 +08:00
fix(agent): retry failed node config applies (#40)
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user