From f29292dd8158b6f705155b1668dd0faded117951 Mon Sep 17 00:00:00 2001 From: ryan Date: Tue, 2 Jun 2026 21:08:04 +0800 Subject: [PATCH] =?UTF-8?q?[=E4=BC=98=E5=8C=96]=20=E6=B7=BB=E5=8A=A0?= =?UTF-8?q?=E8=BF=9B=E7=A8=8B=E7=AE=A1=E7=90=86=E5=8A=9F=E8=83=BD=EF=BC=8C?= =?UTF-8?q?=E6=94=AF=E6=8C=81PID=E6=96=87=E4=BB=B6=E5=A4=84=E7=90=86?= =?UTF-8?q?=E5=92=8C=E5=AD=A4=E5=84=BF=E8=BF=9B=E7=A8=8B=E6=B8=85=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- openflare_relay/internal/frps/manager.go | 30 ++++++++ openflare_relay/internal/frps/manager_test.go | 50 +++++++++++++ openflared/internal/frpc/manager.go | 35 +++++++++- openflared/internal/frpc/manager_test.go | 70 +++++++++++++++++++ 4 files changed, 184 insertions(+), 1 deletion(-) diff --git a/openflare_relay/internal/frps/manager.go b/openflare_relay/internal/frps/manager.go index 6a90c114..650b487e 100644 --- a/openflare_relay/internal/frps/manager.go +++ b/openflare_relay/internal/frps/manager.go @@ -19,6 +19,7 @@ type Manager struct { frpsPath string dataDir string configPath string + pidPath string agentToken string mu sync.RWMutex @@ -45,6 +46,7 @@ func NewManager(frpsPath string, dataDir string, agentToken string) *Manager { frpsPath: frpsPath, dataDir: dataDir, configPath: filepath.Join(dataDir, "frps.toml"), + pidPath: filepath.Join(dataDir, "frps.pid"), status: "unknown", // 启动阶段尚未获取配置,状态未知;避免首次 heartbeat 误报 frps_unhealthy agentToken: agentToken, } @@ -183,6 +185,8 @@ func (m *Manager) supervise(generation uint64) { return } + ensureNoOrphanProcess(m.pidPath) + cmd := exec.Command(m.frpsPath, "-c", m.configPath) cmd.Stdout = os.Stdout cmd.Stderr = os.Stderr @@ -204,6 +208,8 @@ func (m *Manager) supervise(generation uint64) { continue } + _ = os.WriteFile(m.pidPath, []byte(fmt.Sprintf("%d", cmd.Process.Pid)), 0644) + m.cmd = cmd m.status = "healthy" m.lastError = "" @@ -211,6 +217,7 @@ func (m *Manager) supervise(generation uint64) { startedAt := time.Now() waitErr := cmd.Wait() + _ = os.Remove(m.pidPath) m.mu.Lock() if m.cmd == cmd { @@ -272,5 +279,28 @@ func (m *Manager) Stop() { _ = m.cmd.Process.Kill() m.cmd = nil } + _ = os.Remove(m.pidPath) m.status = "unhealthy" } + +func ensureNoOrphanProcess(pidPath string) { + data, err := os.ReadFile(pidPath) + if err != nil { + return + } + var pid int + if _, err := fmt.Sscanf(string(data), "%d", &pid); err != nil { + return + } + if pid <= 0 { + return + } + process, err := os.FindProcess(pid) + if err == nil && process != nil { + slog.Warn("attempting to kill potentially orphan process", "pid", pid, "pid_path", pidPath) + _ = process.Kill() + // Wait a little bit to ensure the OS has reclaimed ports + time.Sleep(500 * time.Millisecond) + } + _ = os.Remove(pidPath) +} diff --git a/openflare_relay/internal/frps/manager_test.go b/openflare_relay/internal/frps/manager_test.go index a2f90621..d4092e07 100644 --- a/openflare_relay/internal/frps/manager_test.go +++ b/openflare_relay/internal/frps/manager_test.go @@ -3,6 +3,7 @@ package frps import ( "fmt" "os" + "os/exec" "path/filepath" "strings" "sync/atomic" @@ -63,6 +64,21 @@ func assertStatusEventually(t *testing.T, m *Manager, expectedStatus string, tim t.Fatalf("expected status eventually %s, got %s (err: %s)", expectedStatus, rt.Status, rt.LastError) } +func assertCommandExitedEventually(t *testing.T, cmd *exec.Cmd, timeout time.Duration) { + t.Helper() + + done := make(chan error, 1) + go func() { + done <- cmd.Wait() + }() + + select { + case <-time.After(timeout): + t.Fatalf("expected process pid=%d to exit within %s", cmd.Process.Pid, timeout) + case <-done: + } +} + func TestStartProcessSuccess(t *testing.T) { scriptPath, dir := setupDummyScript(t) writeControl(t, dir, 0, 5) // exit code 0, sleep 5s @@ -293,3 +309,37 @@ func TestSupervisorGenerationInterrupt(t *testing.T) { t.Error("expected first process to be killed") } } + +func TestUpdateConfigKillsOrphanProcessBeforeRestart(t *testing.T) { + scriptPath, dir := setupDummyScript(t) + writeControl(t, dir, 0, 5) + + m := NewManager(scriptPath, dir, "agent-token") + defer m.Stop() + + orphan := exec.Command("sh", "-c", "sleep 30") + if err := orphan.Start(); err != nil { + t.Fatalf("failed to start orphan process: %v", err) + } + t.Cleanup(func() { + if orphan.Process != nil { + _ = orphan.Process.Kill() + } + }) + + if err := os.WriteFile(m.pidPath, []byte(fmt.Sprintf("%d", orphan.Process.Pid)), 0o644); err != nil { + t.Fatalf("failed to seed orphan pid file: %v", err) + } + + cfg := &service.RelayConfig{ + BindPort: 7000, + VhostHTTPPort: 8080, + AuthToken: "test-auth", + WebServerEnabled: false, + } + + m.UpdateConfig(cfg) + + assertCommandExitedEventually(t, orphan, 2*time.Second) + assertStatusEventually(t, m, "healthy", 2*time.Second) +} diff --git a/openflared/internal/frpc/manager.go b/openflared/internal/frpc/manager.go index de84575f..b6528931 100644 --- a/openflared/internal/frpc/manager.go +++ b/openflared/internal/frpc/manager.go @@ -129,6 +129,8 @@ func (m *Manager) UpdateConfig(ctx context.Context, newConfig *service.FlaredTun if _, ok := activeRelays[relayID]; !ok { slog.Info("stopping obsolete frpc process", "relay_id", relayID) proc.Cancel() + pidPath := filepath.Join(m.cfg.DataDir, fmt.Sprintf("frpc_%s.pid", relayID)) + _ = os.Remove(pidPath) delete(m.processes, relayID) } } @@ -142,8 +144,10 @@ func (m *Manager) UpdateConfig(ctx context.Context, newConfig *service.FlaredTun } func (m *Manager) restartProcess(ctx context.Context, relayID string, configPath string) { + pidPath := filepath.Join(m.cfg.DataDir, fmt.Sprintf("frpc_%s.pid", relayID)) if proc, ok := m.processes[relayID]; ok { proc.Cancel() + _ = os.Remove(pidPath) } procCtx, cancel := context.WithCancel(context.Background()) @@ -167,6 +171,8 @@ func (m *Manager) restartProcess(ctx context.Context, relayID string, configPath } m.mu.Unlock() + ensureNoOrphanProcess(pidPath) + cmd := exec.CommandContext(procCtx, m.cfg.FrpcPath, "-c", configPath) m.mu.Lock() @@ -175,7 +181,12 @@ func (m *Manager) restartProcess(ctx context.Context, relayID string, configPath m.mu.Unlock() startedAt := time.Now() - err := cmd.Run() + err := cmd.Start() + if err == nil { + _ = os.WriteFile(pidPath, []byte(fmt.Sprintf("%d", cmd.Process.Pid)), 0o644) + err = cmd.Wait() + } + _ = os.Remove(pidPath) m.mu.Lock() if procCtx.Err() != nil { @@ -298,3 +309,25 @@ func (m *Manager) LoadState() error { m.mu.Unlock() return nil } + +func ensureNoOrphanProcess(pidPath string) { + data, err := os.ReadFile(pidPath) + if err != nil { + return + } + var pid int + if _, err := fmt.Sscanf(string(data), "%d", &pid); err != nil { + return + } + if pid <= 0 { + return + } + process, err := os.FindProcess(pid) + if err == nil && process != nil { + slog.Warn("attempting to kill potentially orphan process", "pid", pid, "pid_path", pidPath) + _ = process.Kill() + // Wait a little bit to ensure the OS has reclaimed ports + time.Sleep(500 * time.Millisecond) + } + _ = os.Remove(pidPath) +} diff --git a/openflared/internal/frpc/manager_test.go b/openflared/internal/frpc/manager_test.go index a5c7608e..5e0d0bb1 100644 --- a/openflared/internal/frpc/manager_test.go +++ b/openflared/internal/frpc/manager_test.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "os" + "os/exec" "path/filepath" "strings" "testing" @@ -76,6 +77,21 @@ func assertStatusEventually(t *testing.T, m *Manager, relayID string, expectedSt t.Fatalf("expected status eventually %s, got %s (err: %s)", expectedStatus, got, errStr) } +func assertCommandExitedEventually(t *testing.T, cmd *exec.Cmd, timeout time.Duration) { + t.Helper() + + done := make(chan error, 1) + go func() { + done <- cmd.Wait() + }() + + select { + case <-time.After(timeout): + t.Fatalf("expected process pid=%d to exit within %s", cmd.Process.Pid, timeout) + case <-done: + } +} + func TestStartProcessSuccess(t *testing.T) { scriptPath, dir := setupDummyScript(t) writeControl(t, dir, 0, 5) // exit code 0, sleep 5s @@ -265,3 +281,57 @@ func TestBackoffReset(t *testing.T) { m.mu.RUnlock() proc.Cancel() } + +func TestUpdateConfigKillsOrphanProcessBeforeRestart(t *testing.T) { + scriptPath, dir := setupDummyScript(t) + writeControl(t, dir, 0, 5) + + cfg := &config.Config{ + ServerURL: "http://localhost:8080", + TunnelToken: "test-token", + FrpcPath: scriptPath, + DataDir: dir, + StatePath: filepath.Join(dir, "flared-state.json"), + } + + m := NewManager(cfg) + + orphan := exec.Command("sh", "-c", "sleep 30") + if err := orphan.Start(); err != nil { + t.Fatalf("failed to start orphan process: %v", err) + } + t.Cleanup(func() { + if orphan.Process != nil { + _ = orphan.Process.Kill() + } + }) + + pidPath := filepath.Join(dir, "frpc_relay-1.pid") + if err := os.WriteFile(pidPath, []byte(fmt.Sprintf("%d", orphan.Process.Pid)), 0o644); err != nil { + t.Fatalf("failed to seed orphan pid file: %v", err) + } + + newConfig := &service.FlaredTunnelConfigResponse{ + Version: "1", + Checksum: "sum1", + Relays: []service.FlaredRelayInfo{ + { + RelayNodeID: "relay-1", + Address: "127.0.0.1:7000", + AuthToken: "auth-1", + }, + }, + } + + if err := m.UpdateConfig(context.Background(), newConfig); err != nil { + t.Fatalf("failed to UpdateConfig: %v", err) + } + + assertCommandExitedEventually(t, orphan, 2*time.Second) + assertStatusEventually(t, m, "relay-1", "running", 4*time.Second) + + m.mu.RLock() + proc := m.processes["relay-1"] + m.mu.RUnlock() + proc.Cancel() +}