From 7252fb6285303c2d195ed65493a486344c7f9637 Mon Sep 17 00:00:00 2001 From: ryan Date: Tue, 2 Jun 2026 17:12:34 +0800 Subject: [PATCH] =?UTF-8?q?[=E4=BC=98=E5=8C=96]=20=E9=87=8D=E6=9E=84?= =?UTF-8?q?=E8=BF=9B=E7=A8=8B=E7=AE=A1=E7=90=86=E9=80=BB=E8=BE=91=EF=BC=8C?= =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E8=87=AA=E5=8A=A8=E9=87=8D=E5=90=AF=E5=92=8C?= =?UTF-8?q?=E9=80=80=E9=81=BF=E6=9C=BA=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- openflare_relay/internal/frps/manager.go | 158 ++++++---- openflare_relay/internal/frps/manager_test.go | 295 ++++++++++++++++++ 2 files changed, 393 insertions(+), 60 deletions(-) create mode 100644 openflare_relay/internal/frps/manager_test.go diff --git a/openflare_relay/internal/frps/manager.go b/openflare_relay/internal/frps/manager.go index f4f7672d..6a90c114 100644 --- a/openflare_relay/internal/frps/manager.go +++ b/openflare_relay/internal/frps/manager.go @@ -102,19 +102,32 @@ func (m *Manager) UpdateConfig(cfg *service.RelayConfig) { m.activeConfig.WebServerEnabled == cfg.WebServerEnabled { if m.cmd == nil && !m.stopping { slog.Warn("frps config unchanged but process is not running, restarting") - if err := m.restartProcess(); err != nil { + m.stopping = false + m.generation++ + generation := m.generation + if err := m.renderConfig(cfg); err != nil { + slog.Error("failed to render frps config", "error", err) m.status = "unhealthy" m.lastError = err.Error() - slog.Error("failed to restart frps with unchanged config", "error", err) + return } + go m.supervise(generation) } return } m.activeConfig = cfg m.stopping = false + m.generation++ + generation := m.generation slog.Info("relay config updated, reloading frps") + if m.cmd != nil && m.cmd.Process != nil { + slog.Debug("stopping existing frps process") + _ = m.cmd.Process.Kill() + m.cmd = nil + } + if err := m.renderConfig(cfg); err != nil { slog.Error("failed to render frps config", "error", err) m.status = "unhealthy" @@ -122,14 +135,7 @@ func (m *Manager) UpdateConfig(cfg *service.RelayConfig) { return } - if err := m.restartProcess(); err != nil { - slog.Error("failed to restart frps", "error", err) - m.status = "unhealthy" - m.lastError = err.Error() - } else { - m.status = "healthy" - m.lastError = "" - } + go m.supervise(generation) } func (m *Manager) renderConfig(cfg *service.RelayConfig) error { @@ -166,63 +172,95 @@ func (m *Manager) renderConfig(cfg *service.RelayConfig) error { return os.WriteFile(m.configPath, buf.Bytes(), 0644) } -func (m *Manager) restartProcess() error { - m.generation++ - generation := m.generation - if m.cmd != nil && m.cmd.Process != nil { - slog.Debug("stopping existing frps process") - _ = m.cmd.Process.Kill() - m.cmd = nil - } - return m.startProcessLocked(generation) -} +func (m *Manager) supervise(generation uint64) { + backoff := 1 * time.Second + const maxBackoff = 60 * time.Second -func (m *Manager) startProcessLocked(generation uint64) error { - cmd := exec.Command(m.frpsPath, "-c", m.configPath) - cmd.Stdout = os.Stdout - cmd.Stderr = os.Stderr - - if err := cmd.Start(); err != nil { - return err - } - - m.cmd = cmd - m.status = "healthy" - m.lastError = "" - - go func(c *exec.Cmd) { - err := c.Wait() - slog.Warn("frps process exited", "error", err) + for { m.mu.Lock() - if m.cmd == c { + if m.stopping || m.generation != generation { + m.mu.Unlock() + return + } + + cmd := exec.Command(m.frpsPath, "-c", m.configPath) + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + + err := cmd.Start() + if err != nil { + m.status = "unhealthy" + m.lastError = fmt.Sprintf("failed to start: %v", err) + slog.Error("failed to start frps", "error", err, "generation", generation) + m.mu.Unlock() + + if !m.sleepOrInterrupt(generation, backoff) { + return + } + backoff = backoff * 2 + if backoff > maxBackoff { + backoff = maxBackoff + } + continue + } + + m.cmd = cmd + m.status = "healthy" + m.lastError = "" + m.mu.Unlock() + + startedAt := time.Now() + waitErr := cmd.Wait() + + m.mu.Lock() + if m.cmd == cmd { m.cmd = nil m.status = "unhealthy" - if err != nil { - m.lastError = err.Error() + if waitErr != nil { + m.lastError = fmt.Sprintf("exited with error: %v", waitErr) } else { - m.lastError = "frps process exited" + m.lastError = "exited unexpectedly" + } + slog.Warn("frps process exited unexpectedly", "error", waitErr, "generation", generation) + } + shouldContinue := !m.stopping && m.generation == generation + m.mu.Unlock() + + if !shouldContinue { + return + } + + if time.Since(startedAt) >= 10*time.Second { + backoff = 1 * time.Second + } + + if !m.sleepOrInterrupt(generation, backoff) { + return + } + backoff = backoff * 2 + if backoff > maxBackoff { + backoff = maxBackoff + } + } +} + +func (m *Manager) sleepOrInterrupt(generation uint64, d time.Duration) bool { + ticker := time.NewTicker(100 * time.Millisecond) + defer ticker.Stop() + + deadline := time.Now().Add(d) + for time.Now().Before(deadline) { + select { + case <-ticker.C: + m.mu.RLock() + interrupted := m.stopping || m.generation != generation + m.mu.RUnlock() + if interrupted { + return false } } - shouldRestart := !m.stopping && m.generation == generation - m.mu.Unlock() - if !shouldRestart { - return - } - time.Sleep(2 * time.Second) - m.mu.Lock() - defer m.mu.Unlock() - if m.stopping || m.generation != generation { - return - } - slog.Warn("restarting frps after unexpected exit") - if err := m.startProcessLocked(generation); err != nil { - m.status = "unhealthy" - m.lastError = err.Error() - slog.Error("failed to auto restart frps", "error", err) - } - }(cmd) - - return nil + } + return true } func (m *Manager) Stop() { diff --git a/openflare_relay/internal/frps/manager_test.go b/openflare_relay/internal/frps/manager_test.go new file mode 100644 index 00000000..a2f90621 --- /dev/null +++ b/openflare_relay/internal/frps/manager_test.go @@ -0,0 +1,295 @@ +package frps + +import ( + "fmt" + "os" + "path/filepath" + "strings" + "sync/atomic" + "testing" + "time" + + "openflare/service" +) + +// Helper to write control file for the dummy script +func writeControl(t *testing.T, dir string, exitCode int, delaySeconds int) { + controlPath := filepath.Join(dir, "control.txt") + content := fmt.Sprintf("%d %d\n", exitCode, delaySeconds) + err := os.WriteFile(controlPath, []byte(content), 0644) + if err != nil { + t.Fatalf("failed to write control file: %v", err) + } +} + +// Setup a dummy executable script that reads control.txt to decide exit code and sleep duration +func setupDummyScript(t *testing.T) (string, string) { + dir := t.TempDir() + scriptPath := filepath.Join(dir, "dummy_frps") + + // On macOS/Linux, we write a shell script + scriptContent := fmt.Sprintf(`#!/bin/sh +control_file="%s/control.txt" +EXIT_CODE=0 +DELAY=0 +if [ -f "$control_file" ]; then + read -r EXIT_CODE DELAY < "$control_file" +fi +if [ -n "$DELAY" ] && [ "$DELAY" -gt 0 ] 2>/dev/null; then + sleep "$DELAY" +fi +exit "${EXIT_CODE:-0}" +`, dir) + + err := os.WriteFile(scriptPath, []byte(scriptContent), 0755) + if err != nil { + t.Fatalf("failed to write dummy script: %v", err) + } + + return scriptPath, dir +} + +// Helper to poll for status to eliminate timing flakiness in tests +func assertStatusEventually(t *testing.T, m *Manager, expectedStatus string, timeout time.Duration) { + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + rt := m.GetRuntimeStatus() + if rt.Status == expectedStatus { + return + } + time.Sleep(50 * time.Millisecond) + } + rt := m.GetRuntimeStatus() + t.Fatalf("expected status eventually %s, got %s (err: %s)", expectedStatus, rt.Status, rt.LastError) +} + +func TestStartProcessSuccess(t *testing.T) { + scriptPath, dir := setupDummyScript(t) + writeControl(t, dir, 0, 5) // exit code 0, sleep 5s + + m := NewManager(scriptPath, dir, "agent-token") + defer m.Stop() + + cfg := &service.RelayConfig{ + BindPort: 7000, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + + m.UpdateConfig(cfg) + + assertStatusEventually(t, m, "healthy", 2*time.Second) + + rt := m.GetRuntimeStatus() + if !rt.ProcessAlive { + t.Error("expected process to be alive") + } +} + +func TestStartProcessFailureAndBackoff(t *testing.T) { + dir := t.TempDir() + invalidScriptPath := filepath.Join(dir, "non_existent_frps") + + m := NewManager(invalidScriptPath, dir, "agent-token") + defer m.Stop() + + cfg := &service.RelayConfig{ + BindPort: 7000, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + + m.UpdateConfig(cfg) + + assertStatusEventually(t, m, "unhealthy", 2*time.Second) + + rt := m.GetRuntimeStatus() + if !strings.Contains(rt.LastError, "failed to start") { + t.Errorf("expected error message containing 'failed to start', got %s", rt.LastError) + } + + // Correct the path to dummy script + scriptPath, _ := setupDummyScript(t) + writeControl(t, filepath.Dir(scriptPath), 0, 5) + + m.mu.Lock() + m.frpsPath = scriptPath + m.mu.Unlock() + + // Wait for backoff retry (1s backoff) + assertStatusEventually(t, m, "healthy", 3*time.Second) + + rt = m.GetRuntimeStatus() + if !rt.ProcessAlive { + t.Error("expected process to be alive now") + } +} + +func TestUnexpectedExitAndAutorestart(t *testing.T) { + scriptPath, dir := setupDummyScript(t) + // Start with immediate exit code 1 + writeControl(t, dir, 1, 0) + + m := NewManager(scriptPath, dir, "agent-token") + defer m.Stop() + + cfg := &service.RelayConfig{ + BindPort: 7000, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + + m.UpdateConfig(cfg) + + assertStatusEventually(t, m, "unhealthy", 2*time.Second) + + rt := m.GetRuntimeStatus() + if !strings.Contains(rt.LastError, "exited with error") { + t.Errorf("expected exit error, got %s", rt.LastError) + } + + // Change control to be healthy (runs for 5s, exit 0) + writeControl(t, dir, 0, 5) + + // Wait for the retry to fire (backoff was 1s) + assertStatusEventually(t, m, "healthy", 3*time.Second) +} + +func TestBackoffReset(t *testing.T) { + scriptPath, dir := setupDummyScript(t) + // Rapid exit to increase backoff + writeControl(t, dir, 1, 0) + + m := NewManager(scriptPath, dir, "agent-token") + defer m.Stop() + + cfg := &service.RelayConfig{ + BindPort: 7000, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + + m.UpdateConfig(cfg) + + // Crashed once, backoff is 2s + assertStatusEventually(t, m, "unhealthy", 2*time.Second) + + // Now make it run successfully for 11 seconds (exit code 0, sleep 11s) + writeControl(t, dir, 0, 11) + + // Wait for next retry to start running + assertStatusEventually(t, m, "healthy", 4*time.Second) + + // Wait for process to run for 10.5 seconds to trigger backoff reset + time.Sleep(10500 * time.Millisecond) + + // Now make it crash again (exit code 1, sleep 0s) + writeControl(t, dir, 1, 0) + + // Wait for it to finish and crash + assertStatusEventually(t, m, "unhealthy", 3*time.Second) + + // It crashed. Since it ran for > 10s, backoff should have been reset to 1s. + // We make it healthy again (exit code 0, sleep 5) + writeControl(t, dir, 0, 5) + + // Wait 1.5 seconds. If backoff was reset to 1s, it should be healthy now. + assertStatusEventually(t, m, "healthy", 2*time.Second) +} + +func TestImmediateRestartOnSameConfigDeadProcess(t *testing.T) { + scriptPath, dir := setupDummyScript(t) + // Crashes immediately + writeControl(t, dir, 1, 0) + + m := NewManager(scriptPath, dir, "agent-token") + defer m.Stop() + + cfg := &service.RelayConfig{ + BindPort: 7000, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + + m.UpdateConfig(cfg) + + // Let it crash + assertStatusEventually(t, m, "unhealthy", 2*time.Second) + + // Make it start successfully + writeControl(t, dir, 0, 5) + + // Send same config block to trigger immediate restart bypass of backoff sleep + m.UpdateConfig(cfg) + + // Check if it started immediately + assertStatusEventually(t, m, "healthy", 2*time.Second) +} + +func TestSupervisorGenerationInterrupt(t *testing.T) { + scriptPath, dir := setupDummyScript(t) + writeControl(t, dir, 0, 10) + + m := NewManager(scriptPath, dir, "agent-token") + defer m.Stop() + + cfg := &service.RelayConfig{ + BindPort: 7000, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + + m.UpdateConfig(cfg) + + assertStatusEventually(t, m, "healthy", 2*time.Second) + + m.mu.Lock() + gen1 := m.generation + cmd1 := m.cmd + m.mu.Unlock() + + if cmd1 == nil { + t.Fatal("expected active process") + } + + // Update configuration with new bind port to trigger new generation + cfg2 := &service.RelayConfig{ + BindPort: 7001, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + m.UpdateConfig(cfg2) + + assertStatusEventually(t, m, "healthy", 2*time.Second) + + m.mu.Lock() + gen2 := m.generation + cmd2 := m.cmd + m.mu.Unlock() + + if gen2 <= gen1 { + t.Errorf("expected generation incremented, got gen1=%d gen2=%d", gen1, gen2) + } + if cmd2 == cmd1 { + t.Error("expected old process killed and new command started") + } + + // Verify old process is actually killed + var cmd1Finished int32 + go func() { + _ = cmd1.Wait() + atomic.StoreInt32(&cmd1Finished, 1) + }() + + time.Sleep(200 * time.Millisecond) + if atomic.LoadInt32(&cmd1Finished) != 1 { + t.Error("expected first process to be killed") + } +}