[优化] 添加进程管理功能,支持PID文件处理和孤儿进程清理

This commit is contained in:
ryan
2026-06-02 21:08:04 +08:00
parent 4566fc1f53
commit f29292dd81
4 changed files with 184 additions and 1 deletions
+30
View File
@@ -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)
}
@@ -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)
}
+34 -1
View File
@@ -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)
}
+70
View File
@@ -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()
}