[优化] 重构进程管理逻辑,添加自动重启和退避机制

This commit is contained in:
ryan
2026-06-02 17:12:34 +08:00
parent 2220e45989
commit 7252fb6285
2 changed files with 393 additions and 60 deletions
+98 -60
View File
@@ -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() {
@@ -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")
}
}