mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-05 07:26:36 +08:00
[功能] 添加节点健康事件清理功能,优化节点观测数据管理
This commit is contained in:
@@ -105,13 +105,27 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ
|
||||
}
|
||||
snapshot.CurrentVersion = target.Version
|
||||
snapshot.CurrentChecksum = target.Checksum
|
||||
clearBlockedTarget(snapshot)
|
||||
snapshot.LastError = ""
|
||||
slog.Debug("sync finished without changes", "mode", mode, "version", target.Version)
|
||||
return s.stateStore.Save(snapshot)
|
||||
}
|
||||
if isBlockedTarget(snapshot, target.Version, target.Checksum) {
|
||||
slog.Warn("skipping blocked config version after previous failed apply", "mode", mode, "version", target.Version, "checksum", target.Checksum)
|
||||
if startup {
|
||||
if err = s.ensureRuntimeForCurrentConfig(ctx, mode, snapshot, currentChecksum); err != nil {
|
||||
return err
|
||||
}
|
||||
return s.stateStore.Save(snapshot)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if hasBlockedTarget(snapshot) {
|
||||
clearBlockedTarget(snapshot)
|
||||
}
|
||||
if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup {
|
||||
slog.Debug("skipping config fetch because state already records target version/checksum", "version", target.Version, "checksum", target.Checksum)
|
||||
return nil
|
||||
return s.stateStore.Save(snapshot)
|
||||
}
|
||||
|
||||
config, err := s.client.GetActiveConfig(ctx)
|
||||
@@ -139,6 +153,7 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
||||
}
|
||||
snapshot.CurrentVersion = config.Version
|
||||
snapshot.CurrentChecksum = config.Checksum
|
||||
clearBlockedTarget(snapshot)
|
||||
snapshot.LastError = ""
|
||||
slog.Debug("sync finished without changes", "mode", mode, "version", config.Version)
|
||||
return s.stateStore.Save(snapshot)
|
||||
@@ -146,9 +161,22 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
||||
if target != nil && (target.Version != config.Version || target.Checksum != config.Checksum) {
|
||||
slog.Warn("active config changed between heartbeat and fetch", "heartbeat_version", target.Version, "heartbeat_checksum", target.Checksum, "fetched_version", config.Version, "fetched_checksum", config.Checksum)
|
||||
}
|
||||
if isBlockedTarget(snapshot, config.Version, config.Checksum) {
|
||||
slog.Warn("skipping blocked config after fetch because the same version previously failed", "mode", mode, "version", config.Version, "checksum", config.Checksum)
|
||||
if startup {
|
||||
if err := s.ensureRuntimeForCurrentConfig(ctx, mode, snapshot, currentChecksum); err != nil {
|
||||
return err
|
||||
}
|
||||
return s.stateStore.Save(snapshot)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if hasBlockedTarget(snapshot) {
|
||||
clearBlockedTarget(snapshot)
|
||||
}
|
||||
if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum && !startup {
|
||||
slog.Debug("skipping apply because state already records target version/checksum", "version", config.Version, "checksum", config.Checksum)
|
||||
return nil
|
||||
return s.stateStore.Save(snapshot)
|
||||
}
|
||||
routeConfig := config.RouteConfig
|
||||
if routeConfig == "" {
|
||||
@@ -172,6 +200,7 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
||||
slog.Info("openresty config applied successfully", "mode", mode, "version", config.Version)
|
||||
snapshot.CurrentVersion = config.Version
|
||||
snapshot.CurrentChecksum = config.Checksum
|
||||
clearBlockedTarget(snapshot)
|
||||
snapshot.LastError = ""
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||
snapshot.OpenrestyMessage = ""
|
||||
@@ -184,6 +213,7 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
||||
message = "apply rolled back to previous config"
|
||||
}
|
||||
slog.Warn("openresty config apply rolled back", "mode", mode, "version", config.Version, "message", message)
|
||||
markBlockedTarget(snapshot, config.Version, config.Checksum, message)
|
||||
snapshot.LastError = message
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||
snapshot.OpenrestyMessage = message
|
||||
@@ -193,6 +223,7 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
||||
message = "openresty apply failed"
|
||||
}
|
||||
slog.Error("apply openresty config failed", "mode", mode, "version", config.Version, "message", message)
|
||||
markBlockedTarget(snapshot, config.Version, config.Checksum, message)
|
||||
snapshot.LastError = message
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
|
||||
snapshot.OpenrestyMessage = message
|
||||
@@ -230,6 +261,56 @@ func outcomeError(version string, message string) error {
|
||||
return fmt.Errorf("apply version %s failed: %s", version, trimmed)
|
||||
}
|
||||
|
||||
func (s *Service) ensureRuntimeForCurrentConfig(ctx context.Context, mode string, snapshot *state.Snapshot, currentChecksum string) error {
|
||||
if strings.TrimSpace(currentChecksum) == "" {
|
||||
slog.Warn("blocked config cannot be retried and no local checksum is available for runtime recovery", "mode", mode, "blocked_version", snapshot.BlockedVersion)
|
||||
return nil
|
||||
}
|
||||
slog.Info("ensuring runtime with current local config while active target remains blocked", "mode", mode, "current_version", snapshot.CurrentVersion, "current_checksum", currentChecksum, "blocked_version", snapshot.BlockedVersion)
|
||||
if err := s.nginxManager.EnsureRuntime(ctx, true); err != nil {
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
|
||||
snapshot.OpenrestyMessage = err.Error()
|
||||
_ = s.stateStore.Save(snapshot)
|
||||
return err
|
||||
}
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||
if strings.TrimSpace(snapshot.OpenrestyMessage) == strings.TrimSpace(snapshot.BlockedReason) {
|
||||
snapshot.OpenrestyMessage = ""
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func markBlockedTarget(snapshot *state.Snapshot, version string, checksum string, reason string) {
|
||||
if snapshot == nil {
|
||||
return
|
||||
}
|
||||
snapshot.BlockedVersion = strings.TrimSpace(version)
|
||||
snapshot.BlockedChecksum = strings.TrimSpace(checksum)
|
||||
snapshot.BlockedReason = strings.TrimSpace(reason)
|
||||
}
|
||||
|
||||
func clearBlockedTarget(snapshot *state.Snapshot) {
|
||||
if snapshot == nil {
|
||||
return
|
||||
}
|
||||
snapshot.BlockedVersion = ""
|
||||
snapshot.BlockedChecksum = ""
|
||||
snapshot.BlockedReason = ""
|
||||
}
|
||||
|
||||
func hasBlockedTarget(snapshot *state.Snapshot) bool {
|
||||
return snapshot != nil && (strings.TrimSpace(snapshot.BlockedVersion) != "" || strings.TrimSpace(snapshot.BlockedChecksum) != "")
|
||||
}
|
||||
|
||||
func isBlockedTarget(snapshot *state.Snapshot, version string, checksum string) bool {
|
||||
if snapshot == nil {
|
||||
return false
|
||||
}
|
||||
return strings.TrimSpace(snapshot.BlockedVersion) == strings.TrimSpace(version) &&
|
||||
strings.TrimSpace(snapshot.BlockedChecksum) == strings.TrimSpace(checksum) &&
|
||||
(strings.TrimSpace(version) != "" || strings.TrimSpace(checksum) != "")
|
||||
}
|
||||
|
||||
func checksumString(content string) string {
|
||||
sum := sha256.Sum256([]byte(content))
|
||||
return hex.EncodeToString(sum[:])
|
||||
|
||||
@@ -203,6 +203,9 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) {
|
||||
if snapshot.CurrentVersion != "20260309-001" {
|
||||
t.Fatal("expected failed sync not to overwrite current version")
|
||||
}
|
||||
if snapshot.BlockedVersion != "20260309-002" || snapshot.BlockedChecksum != "checksum-2" {
|
||||
t.Fatalf("expected failed target version to be blocked, got %+v", snapshot)
|
||||
}
|
||||
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusUnhealthy {
|
||||
t.Fatalf("expected unhealthy openresty status, got %q", snapshot.OpenrestyStatus)
|
||||
}
|
||||
@@ -267,6 +270,9 @@ func TestSyncOnceReportsWarningWhenRollbackKeepsOpenrestyHealthy(t *testing.T) {
|
||||
if snapshot.CurrentVersion != "20260309-001" || snapshot.CurrentChecksum != "checksum-1" {
|
||||
t.Fatal("expected warning apply to keep previous version state")
|
||||
}
|
||||
if snapshot.BlockedVersion != "20260309-002" || snapshot.BlockedChecksum != "checksum-2" {
|
||||
t.Fatalf("expected rolled-back target version to be blocked, got %+v", snapshot)
|
||||
}
|
||||
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusHealthy {
|
||||
t.Fatalf("expected healthy openresty after rollback, got %q", snapshot.OpenrestyStatus)
|
||||
}
|
||||
@@ -368,6 +374,165 @@ func TestSyncOnStartupRecordsRuntimeFailure(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnceSkipsPreviouslyBlockedVersion(t *testing.T) {
|
||||
client := &fakeClient{
|
||||
config: protocol.ActiveConfigResponse{
|
||||
Version: "20260309-006",
|
||||
Checksum: "checksum-6",
|
||||
MainConfig: "worker_processes 6;",
|
||||
RouteConfig: "server { listen 86; }",
|
||||
RenderedConfig: "server { listen 86; }",
|
||||
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: "20260309-005",
|
||||
CurrentChecksum: "checksum-5",
|
||||
BlockedVersion: "20260309-006",
|
||||
BlockedChecksum: "checksum-6",
|
||||
BlockedReason: "apply failed, rolled back to previous config",
|
||||
LastError: "apply failed, rolled back to previous config",
|
||||
}); err != nil {
|
||||
t.Fatalf("failed to seed state: %v", err)
|
||||
}
|
||||
|
||||
manager := &fakeManager{currentChecksum: "checksum-5"}
|
||||
service := New(client, manager, stateStore)
|
||||
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{
|
||||
Version: "20260309-006",
|
||||
Checksum: "checksum-6",
|
||||
}); err != nil {
|
||||
t.Fatalf("expected blocked version to be skipped, got %v", err)
|
||||
}
|
||||
if client.fetchCalls != 0 {
|
||||
t.Fatalf("expected blocked version to skip fetch, got %d", client.fetchCalls)
|
||||
}
|
||||
if len(manager.applyMainContents) != 0 {
|
||||
t.Fatal("expected blocked version to skip apply")
|
||||
}
|
||||
if len(client.reports) != 0 {
|
||||
t.Fatal("expected blocked version to skip reporting duplicate apply result")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnStartupKeepsBlockedVersionSuppressedUntilNewTargetArrives(t *testing.T) {
|
||||
client := &fakeClient{
|
||||
config: protocol.ActiveConfigResponse{
|
||||
Version: "20260309-007",
|
||||
Checksum: "checksum-7",
|
||||
MainConfig: "worker_processes 7;",
|
||||
RouteConfig: "server { listen 87; }",
|
||||
RenderedConfig: "server { listen 87; }",
|
||||
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: "20260309-005",
|
||||
CurrentChecksum: "checksum-5",
|
||||
BlockedVersion: "20260309-007",
|
||||
BlockedChecksum: "checksum-7",
|
||||
BlockedReason: "apply failed, rolled back to previous config",
|
||||
OpenrestyStatus: protocol.OpenrestyStatusUnhealthy,
|
||||
OpenrestyMessage: "apply failed, rolled back to previous config",
|
||||
LastError: "apply failed, rolled back to previous config",
|
||||
}); err != nil {
|
||||
t.Fatalf("failed to seed state: %v", err)
|
||||
}
|
||||
|
||||
manager := &fakeManager{currentChecksum: "checksum-5"}
|
||||
service := New(client, manager, stateStore)
|
||||
if err = service.SyncOnStartup(context.Background(), &protocol.ActiveConfigMeta{
|
||||
Version: "20260309-007",
|
||||
Checksum: "checksum-7",
|
||||
}); err != nil {
|
||||
t.Fatalf("expected blocked startup target to be skipped, got %v", err)
|
||||
}
|
||||
if len(manager.ensureCalls) != 1 || !manager.ensureCalls[0] {
|
||||
t.Fatal("expected startup skip to ensure runtime with current local config")
|
||||
}
|
||||
if client.fetchCalls != 0 {
|
||||
t.Fatalf("expected blocked startup target to skip fetch, got %d", client.fetchCalls)
|
||||
}
|
||||
if len(client.reports) != 0 {
|
||||
t.Fatal("expected blocked startup target to skip duplicate apply report")
|
||||
}
|
||||
snapshot, err := stateStore.Load()
|
||||
if err != nil {
|
||||
t.Fatalf("failed to load state: %v", err)
|
||||
}
|
||||
if snapshot.BlockedVersion != "20260309-007" || snapshot.BlockedChecksum != "checksum-7" {
|
||||
t.Fatalf("expected blocked target to remain recorded, got %+v", snapshot)
|
||||
}
|
||||
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusHealthy {
|
||||
t.Fatalf("expected startup runtime recovery to mark openresty healthy, got %q", snapshot.OpenrestyStatus)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnceClearsBlockedTargetWhenNewVersionArrives(t *testing.T) {
|
||||
client := &fakeClient{
|
||||
config: protocol.ActiveConfigResponse{
|
||||
Version: "20260309-008",
|
||||
Checksum: "checksum-8",
|
||||
MainConfig: "worker_processes 8;",
|
||||
RouteConfig: "server { listen 88; }",
|
||||
RenderedConfig: "server { listen 88; }",
|
||||
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: "20260309-005",
|
||||
CurrentChecksum: "checksum-5",
|
||||
BlockedVersion: "20260309-007",
|
||||
BlockedChecksum: "checksum-7",
|
||||
BlockedReason: "apply failed, rolled back to previous config",
|
||||
}); err != nil {
|
||||
t.Fatalf("failed to seed state: %v", err)
|
||||
}
|
||||
|
||||
manager := &fakeManager{}
|
||||
service := New(client, manager, stateStore)
|
||||
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{
|
||||
Version: "20260309-008",
|
||||
Checksum: "checksum-8",
|
||||
}); err != nil {
|
||||
t.Fatalf("expected new target version to be applied, got %v", err)
|
||||
}
|
||||
if client.fetchCalls != 1 {
|
||||
t.Fatalf("expected new target to trigger fetch, got %d", client.fetchCalls)
|
||||
}
|
||||
if len(manager.applyMainContents) != 1 {
|
||||
t.Fatal("expected new target to trigger apply")
|
||||
}
|
||||
snapshot, err := stateStore.Load()
|
||||
if err != nil {
|
||||
t.Fatalf("failed to load state: %v", err)
|
||||
}
|
||||
if snapshot.BlockedVersion != "" || snapshot.BlockedChecksum != "" {
|
||||
t.Fatalf("expected blocked target to be cleared after new version succeeds, got %+v", snapshot)
|
||||
}
|
||||
if snapshot.CurrentVersion != "20260309-008" || snapshot.CurrentChecksum != "checksum-8" {
|
||||
t.Fatalf("expected current version to move to new target, got %+v", snapshot)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnceSkipsFetchWhenHeartbeatChecksumMatches(t *testing.T) {
|
||||
client := &fakeClient{
|
||||
config: protocol.ActiveConfigResponse{
|
||||
|
||||
Reference in New Issue
Block a user