From 21d1a4d1802e680295a312288f23c884f826d1b8 Mon Sep 17 00:00:00 2001 From: lxh2875931338 Date: Fri, 7 Aug 2026 23:14:13 +0800 Subject: [PATCH] =?UTF-8?q?=E5=85=A8=E6=96=B0=E7=9A=84=E8=BF=9B=E7=A8=8B?= =?UTF-8?q?=E6=8E=A7=E5=88=B6=E5=8A=A8=E4=BD=9C=E7=8A=B6=E6=80=81=E6=9C=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/frp/legacy.go | 6 +- internal/process/attr_linux.go | 1 - internal/process/attr_other.go | 3 +- internal/process/manager.go | 494 +++++++++++++++++++++++++-------- main.go | 6 +- 5 files changed, 379 insertions(+), 131 deletions(-) diff --git a/internal/frp/legacy.go b/internal/frp/legacy.go index ad4108a..646efa6 100644 --- a/internal/frp/legacy.go +++ b/internal/frp/legacy.go @@ -27,7 +27,7 @@ func IsRunning() bool { log.Printf("[WARN] ProcessManager.Status() 失败: %v,降级到 PID 文件", err) return isRunningLegacy() } - return status.State == "running" + return status.Phase == "RUNNING" } return isRunningLegacy() } @@ -77,7 +77,7 @@ func GetStatus() (map[string]interface{}, error) { return nil, err } return map[string]interface{}{ - "state": status.State, + "state": status.Phase, "pid": status.PID, "port": status.Port, }, nil @@ -104,7 +104,7 @@ func Reload() error { return err } - if status.State != "running" { + if status.Phase != "RUNNING" { return pm.Start(context.Background()) } diff --git a/internal/process/attr_linux.go b/internal/process/attr_linux.go index 51472d6..719cab8 100644 --- a/internal/process/attr_linux.go +++ b/internal/process/attr_linux.go @@ -7,7 +7,6 @@ import ( "syscall" ) -// setProcessAttributes 设置 Linux 进程属性 (Setpgid) func setProcessAttributes(cmd *exec.Cmd) { cmd.SysProcAttr = &syscall.SysProcAttr{ Setpgid: true, diff --git a/internal/process/attr_other.go b/internal/process/attr_other.go index 6ca9d61..73cb6dd 100644 --- a/internal/process/attr_other.go +++ b/internal/process/attr_other.go @@ -6,7 +6,6 @@ import ( "os/exec" ) -// setProcessAttributes 其他平台(macOS/FreeBSD 等)不做特殊设置 func setProcessAttributes(cmd *exec.Cmd) { - // 其他平台不需要特殊设置 + // 其他平台不做特殊设置 } diff --git a/internal/process/manager.go b/internal/process/manager.go index 290dd70..20ba30b 100644 --- a/internal/process/manager.go +++ b/internal/process/manager.go @@ -1,6 +1,6 @@ // internal/process/manager.go // frpc-console 进程管理模块 -// 2.6-preview: 端口检测 + 单实例锁定 + 状态自述 +// 2.6-preview: 状态机驱动 + 端口检测 + 单实例锁定 package process @@ -37,16 +37,55 @@ const ( LockRetryInterval = 100 * time.Millisecond APITimeout = 2 * time.Second - // 启动检测参数 - StartupMaxAttempts = 25 // 最多检测 25 次 - StartupRetryDelay = 200 * time.Millisecond // 每次间隔 200ms - // 总超时 = 25 * 200ms = 5 秒 + StartupMaxAttempts = 50 + StartupRetryDelay = 200 * time.Millisecond + StartupTimeout = 10 * time.Second +) + +// ================================================================ +// 进程状态定义 +// ================================================================ + +type ProcessPhase string + +const ( + PhaseUnknown ProcessPhase = "UNKNOWN" + PhaseStarting ProcessPhase = "STARTING" + PhaseRunning ProcessPhase = "RUNNING" + PhaseDegraded ProcessPhase = "DEGRADED" + PhaseFailed ProcessPhase = "FAILED" + PhaseStopped ProcessPhase = "STOPPED" ) // ================================================================ // 数据结构 // ================================================================ +type PortCheckResult struct { + Ready bool + Err error +} + +type ProcessState struct { + Phase ProcessPhase `json:"phase"` + PID int `json:"pid"` + Port int `json:"port"` + StartedAt time.Time `json:"started_at"` + ExitCode int `json:"exit_code,omitempty"` + + Alive bool `json:"alive"` + PortReady bool `json:"port_ready"` + PortError string `json:"port_error,omitempty"` + FRPReady bool `json:"frp_ready"` + Version string `json:"version,omitempty"` + + ExpectedStop bool `json:"expected_stop"` + + LastOutput string `json:"last_output,omitempty"` + LastError string `json:"last_error,omitempty"` + Error string `json:"error,omitempty"` +} + type PortStatus struct { Port int `json:"port"` Occupied bool `json:"occupied"` @@ -55,15 +94,6 @@ type PortStatus struct { ProcessCmd string `json:"process_cmd,omitempty"` } -type ProcessStatus struct { - State string `json:"state"` - PID int `json:"pid"` - Port int `json:"port"` - Uptime string `json:"uptime,omitempty"` - Version string `json:"version,omitempty"` - Error string `json:"error,omitempty"` -} - type FRPCStatus struct { Version string `json:"version"` RunID string `json:"run_id"` @@ -85,8 +115,18 @@ type ProcessManager struct { configPath string frpcBinPath string adminPort int - lockFile *os.File - locked bool + + startTime time.Time + expectedStop bool + exitCode int + exitMu sync.Mutex + + lastOutput string + lastError string + outputMu sync.Mutex + + lockFile *os.File + locked bool } var ( @@ -94,6 +134,10 @@ var ( globalMu sync.Mutex ) +// ================================================================ +// 构造函数 +// ================================================================ + func NewManager(dataDir, configPath, frpcBinPath string) *ProcessManager { return &ProcessManager{ dataDir: dataDir, @@ -234,28 +278,29 @@ func extractPortFromAddrSection(content, section, key string) int { // 端口检测 // ================================================================ -func (pm *ProcessManager) CheckPort() (bool, error) { +func (pm *ProcessManager) CheckPort() PortCheckResult { if pm.adminPort <= 0 { - return false, fmt.Errorf("admin_port 未配置 (当前值: %d)", pm.adminPort) + return PortCheckResult{Ready: false, Err: fmt.Errorf("admin_port 未配置")} } + conn, err := net.DialTimeout("tcp", fmt.Sprintf("127.0.0.1:%d", pm.adminPort), PortCheckTimeout) if err != nil { - return false, nil + return PortCheckResult{Ready: false, Err: err} } conn.Close() - return true, nil + return PortCheckResult{Ready: true, Err: nil} } func (pm *ProcessManager) GetPortStatus() (*PortStatus, error) { status := &PortStatus{Port: pm.adminPort, Occupied: false, PID: 0, IsFRPC: false} - occupied, err := pm.CheckPort() - if err != nil { - return status, err + result := pm.CheckPort() + if result.Err != nil { + return status, result.Err } - status.Occupied = occupied - if !occupied { + if !result.Ready { return status, nil } + status.Occupied = true pid, err := pm.getPIDByPort(pm.adminPort) if err != nil { @@ -442,13 +487,11 @@ func (pm *ProcessManager) killProcess(pid int) error { } func (pm *ProcessManager) cleanupOrphans() { - // 检查端口是否被占用 - occupied, err := pm.CheckPort() - if err != nil || !occupied { + result := pm.CheckPort() + if !result.Ready { return } - // 检查是否有 PID 文件 pid := pm.readPIDFile() if pid > 0 && pm.isProcessAlive(pid) { return @@ -496,39 +539,179 @@ func (pm *ProcessManager) deletePIDFile() error { return err } +// ================================================================ +// stdout/stderr 捕获 +// ================================================================ + +func (pm *ProcessManager) captureOutput(reader io.ReadCloser, isError bool) { + defer reader.Close() + scanner := bufio.NewScanner(reader) + buf := make([]byte, 64*1024) + scanner.Buffer(buf, 256*1024) + + var lines []string + const maxLines = 20 + + for scanner.Scan() { + line := scanner.Text() + pm.outputMu.Lock() + if isError { + pm.lastError = line + } else { + pm.lastOutput = line + } + pm.outputMu.Unlock() + + if isError { + lines = append(lines, line) + if len(lines) > maxLines { + lines = lines[1:] + } + pm.outputMu.Lock() + pm.lastError = strings.Join(lines, "\n") + pm.outputMu.Unlock() + } else { + lines = append(lines, line) + if len(lines) > maxLines { + lines = lines[1:] + } + pm.outputMu.Lock() + pm.lastOutput = strings.Join(lines, "\n") + pm.outputMu.Unlock() + } + } +} + +func (pm *ProcessManager) getLastOutput() string { + pm.outputMu.Lock() + defer pm.outputMu.Unlock() + return pm.lastOutput +} + +func (pm *ProcessManager) getLastError() string { + pm.outputMu.Lock() + defer pm.outputMu.Unlock() + return pm.lastError +} + +// ================================================================ +// FRPReady 检测 +// ================================================================ + +func (pm *ProcessManager) isFRPReady(stdout string) bool { + if strings.Contains(stdout, "start proxy success") || + strings.Contains(stdout, "login to server success") { + return true + } + return false +} + +// ================================================================ +// 状态计算 +// ================================================================ + +func (pm *ProcessManager) computeState(pid int, result PortCheckResult) ProcessState { + alive := pm.isProcessAlive(pid) + portReady := result.Ready + portErr := result.Err + + stdout := pm.getLastOutput() + stderr := pm.getLastError() + frpReady := pm.isFRPReady(stdout) + + // 获取退出码 + pm.exitMu.Lock() + exitCode := pm.exitCode + pm.exitMu.Unlock() + + state := ProcessState{ + PID: pid, + Port: pm.adminPort, + Alive: alive, + PortReady: portReady, + PortError: "", + FRPReady: frpReady, + StartedAt: pm.startTime, + ExpectedStop: pm.expectedStop, + LastOutput: stdout, + LastError: stderr, + ExitCode: exitCode, + } + + if portErr != nil { + state.PortError = portErr.Error() + } + + // ---- 状态判定逻辑 ---- + if !alive { + if pm.expectedStop { + state.Phase = PhaseStopped + return state + } + + // 异常退出 + if exitCode != 0 && exitCode != -1 { + state.Error = fmt.Sprintf("进程异常退出 (exit code: %d)", exitCode) + if stderr != "" { + state.Error += ": " + stderr + } + } else { + state.Error = "进程意外退出" + if stderr != "" { + state.Error += ": " + stderr + } + } + state.Phase = PhaseFailed + return state + } + + // ---- alive == true ---- + if time.Since(pm.startTime) > StartupTimeout { + state.Error = fmt.Sprintf("启动超时 (超过 %v)", StartupTimeout) + if portErr != nil { + state.Error += ": " + portErr.Error() + } + state.Phase = PhaseFailed + return state + } + + if portReady && frpReady { + state.Phase = PhaseRunning + return state + } + + if portReady && !frpReady { + state.Phase = PhaseDegraded + state.Error = "端口已监听,但 frpc 未报告就绪" + return state + } + + if !portReady { + if portErr != nil && strings.Contains(portErr.Error(), "permission denied") { + state.Phase = PhaseDegraded + state.Error = "端口检测权限不足: " + portErr.Error() + return state + } + state.Phase = PhaseStarting + if portErr != nil { + state.Error = "等待端口就绪: " + portErr.Error() + } + return state + } + + state.Phase = PhaseUnknown + return state +} + // ================================================================ // 状态查询 // ================================================================ -func (pm *ProcessManager) Status() (*ProcessStatus, error) { - status := &ProcessStatus{State: "unknown", PID: 0, Port: pm.adminPort} - if pm.adminPort <= 0 { - status.Error = fmt.Sprintf("admin_port 未配置 (当前值: %d)", pm.adminPort) - return status, nil - } - portStatus, err := pm.GetPortStatus() - if err != nil { - status.Error = err.Error() - return status, nil - } - if !portStatus.Occupied { - pm.deletePIDFile() - status.State = "stopped" - return status, nil - } - if !portStatus.IsFRPC { - status.State = "conflict" - status.PID = portStatus.PID - status.Error = fmt.Sprintf("端口 %d 被非 frpc 进程占用 (PID: %d)", pm.adminPort, portStatus.PID) - return status, nil - } - status.State = "running" - status.PID = portStatus.PID - pm.writePIDFile(portStatus.PID) - if info := pm.getFRPCStatus(portStatus.PID); info != nil { - status.Version = info.Version - } - return status, nil +func (pm *ProcessManager) Status() (*ProcessState, error) { + pid := pm.readPIDFile() + result := pm.CheckPort() + state := pm.computeState(pid, result) + return &state, nil } // ================================================================ @@ -567,76 +750,117 @@ func (pm *ProcessManager) getFRPCStatus(pid int) *FRPCStatus { } // ================================================================ -// 操作执行 (Start / Stop / Restart) +// 操作执行 // ================================================================ +// Lock 由平台文件 (lock_*.go) 实现 + func (pm *ProcessManager) Start(ctx context.Context) error { if err := pm.Lock(); err != nil { return fmt.Errorf("获取锁失败: %w", err) } defer pm.Unlock() - // 清理孤儿进程 pm.cleanupOrphans() - return pm.startLocked(ctx) } func (pm *ProcessManager) startLocked(ctx context.Context) error { - portStatus, err := pm.GetPortStatus() - if err != nil { - return fmt.Errorf("检测端口状态失败: %w", err) - } - if portStatus.Occupied { - if portStatus.IsFRPC { - pm.writePIDFile(portStatus.PID) + result := pm.CheckPort() + if result.Ready { + pid := pm.readPIDFile() + if pid > 0 && pm.isProcessAlive(pid) { + log.Printf("[INFO] frpc 已在运行 (PID: %d)", pid) return nil } - return fmt.Errorf("端口 %d 被非 frpc 进程占用 (PID: %d)", pm.adminPort, portStatus.PID) + pm.cleanupOrphans() } + pm.startTime = time.Now() + pm.expectedStop = false + + // 重置退出码 + pm.exitMu.Lock() + pm.exitCode = -1 + pm.exitMu.Unlock() + cmd := exec.CommandContext(ctx, pm.frpcBinPath, "-c", pm.configPath) setProcessAttributes(cmd) - cmd.Stdout = os.Stdout - cmd.Stderr = os.Stderr + + stdoutPipe, err := cmd.StdoutPipe() + if err != nil { + return fmt.Errorf("创建 stdout pipe 失败: %w", err) + } + stderrPipe, err := cmd.StderrPipe() + if err != nil { + return fmt.Errorf("创建 stderr pipe 失败: %w", err) + } + + go pm.captureOutput(stdoutPipe, false) + go pm.captureOutput(stderrPipe, true) if err := cmd.Start(); err != nil { return fmt.Errorf("启动 frpc 失败: %w", err) } - // ⭐ 关键修复:启动 goroutine 回收子进程,避免僵尸 - go func() { - err := cmd.Wait() - if err != nil { - log.Printf("[INFO] frpc 进程 (PID: %d) 已退出: %v", cmd.Process.Pid, err) - } else { - log.Printf("[INFO] frpc 进程 (PID: %d) 正常退出", cmd.Process.Pid) - } - // 进程退出后,如果端口没有被占用,Status 逻辑会清理 PID 文件 - }() - - // 立即写入 PID 文件 if err := pm.writePIDFile(cmd.Process.Pid); err != nil { log.Printf("[WARN] 写入 PID 文件失败: %v", err) } log.Printf("[DEBUG] frpc 进程已启动,PID: %d", cmd.Process.Pid) - // 检测端口是否监听(最多 5 秒) + go func() { + err := cmd.Wait() + var code int + if err != nil { + if exitErr, ok := err.(*exec.ExitError); ok { + code = exitErr.ExitCode() + } else { + code = -1 + } + } else { + code = 0 + } + pm.exitMu.Lock() + pm.exitCode = code + pm.exitMu.Unlock() + log.Printf("[INFO] frpc 进程 (PID: %d) 已退出,退出码: %d", cmd.Process.Pid, code) + + if !pm.expectedStop && code != 0 { + log.Printf("[WARN] frpc 进程异常退出 (exit code: %d)", code) + } + if !pm.expectedStop { + pm.deletePIDFile() + } + }() + for attempt := 0; attempt < StartupMaxAttempts; attempt++ { - occupied, err := pm.CheckPort() - if err == nil && occupied { - log.Printf("[INFO] frpc 启动成功 (PID: %d, 端口: %d),耗时 %dms", - cmd.Process.Pid, pm.adminPort, attempt*int(StartupRetryDelay/time.Millisecond)) - return nil + r := pm.CheckPort() + if r.Ready { + stdout := pm.getLastOutput() + if pm.isFRPReady(stdout) { + log.Printf("[INFO] frpc 启动成功 (PID: %d, 端口: %d),耗时 %dms", + cmd.Process.Pid, pm.adminPort, attempt*int(StartupRetryDelay/time.Millisecond)) + return nil + } + log.Printf("[DEBUG] 端口已就绪,等待 frpc 初始化...") } time.Sleep(StartupRetryDelay) } - // 端口监听超时:杀掉启动的进程并清理 - log.Printf("[WARN] frpc 启动超时 (PID: %d),正在清理...", cmd.Process.Pid) - pm.killProcess(cmd.Process.Pid) + if pm.isProcessAlive(cmd.Process.Pid) { + stderr := pm.getLastError() + log.Printf("[WARN] frpc 启动超时 (PID: %d),当前 stderr: %s", cmd.Process.Pid, stderr) + // 进程还在但端口没起来,尝试获取 frpc admin API 版本信息 + if info := pm.getFRPCStatus(cmd.Process.Pid); info != nil && info.Version != "" { + log.Printf("[INFO] frpc admin API 可访问,版本: %s", info.Version) + // API 可访问说明 frpc 已正常运行,只是端口检测有误 + return nil + } + return fmt.Errorf("frpc 启动超时: 进程存在但端口未就绪") + } + pm.deletePIDFile() - return fmt.Errorf("frpc 启动超时: 端口未监听 (PID: %d)", cmd.Process.Pid) + return fmt.Errorf("frpc 启动失败: 进程已退出") } func (pm *ProcessManager) Stop(ctx context.Context) error { @@ -644,63 +868,56 @@ func (pm *ProcessManager) Stop(ctx context.Context) error { return fmt.Errorf("获取锁失败: %w", err) } defer pm.Unlock() + + pm.expectedStop = true return pm.stopLocked(ctx) } func (pm *ProcessManager) stopLocked(_ context.Context) error { - portStatus, err := pm.GetPortStatus() - if err != nil { - return fmt.Errorf("检测端口状态失败: %w", err) - } - if !portStatus.Occupied { - pm.deletePIDFile() - return nil - } - var pid int - if portStatus.IsFRPC { - pid = portStatus.PID - } else { - return fmt.Errorf("端口 %d 被非 frpc 进程占用, 无法安全停止", pm.adminPort) - } + pid := pm.readPIDFile() if pid <= 0 { - pid = pm.readPIDFile() - if pid <= 0 { - return fmt.Errorf("无法确定 frpc 进程 PID") + if runtime.GOOS == "windows" { + exec.Command("taskkill", "/F", "/IM", "frpc.exe").Run() + } else { + exec.Command("pkill", "-f", "frpc").Run() } + pm.deletePIDFile() + return nil } - // 发送 SIGTERM (优雅停止) - proc, err := os.FindProcess(pid) + if !pm.isProcessAlive(pid) { + pm.deletePIDFile() + return nil + } + + process, err := os.FindProcess(pid) if err != nil { pm.deletePIDFile() return nil } - if err := proc.Signal(syscall.SIGTERM); err != nil { + if err := process.Signal(syscall.SIGTERM); err != nil { pm.deletePIDFile() return nil } - // 等待端口释放 start := time.Now() for time.Since(start) < StopMaxWaitTime { - occupied, _ := pm.CheckPort() - if !occupied { + if !pm.isProcessAlive(pid) { pm.deletePIDFile() return nil } time.Sleep(200 * time.Millisecond) } - // 端口未释放, 强制 kill if runtime.GOOS == "windows" { - exec.Command("taskkill", "/F", "/IM", "frpc.exe").Run() + exec.Command("taskkill", "/F", "/PID", strconv.Itoa(pid)).Run() } else { - proc.Kill() + process.Kill() } time.Sleep(500 * time.Millisecond) - if occupied, _ := pm.CheckPort(); occupied { - return fmt.Errorf("强制停止失败: 端口仍被占用") + if pm.isProcessAlive(pid) { + return fmt.Errorf("强制停止失败: 进程仍存活 (PID: %d)", pid) } pm.deletePIDFile() @@ -712,13 +929,46 @@ func (pm *ProcessManager) Restart(ctx context.Context) error { return fmt.Errorf("获取锁失败: %w", err) } defer pm.Unlock() + if err := pm.stopLocked(ctx); err != nil { return fmt.Errorf("停止失败: %w", err) } - // 等待端口完全释放 time.Sleep(1 * time.Second) if err := pm.startLocked(ctx); err != nil { return fmt.Errorf("启动失败: %w", err) } return nil } + +// ================================================================ +// 健康检查 +// ================================================================ + +func (pm *ProcessManager) HealthCheck() map[string]interface{} { + result := make(map[string]interface{}) + result["admin_port"] = pm.adminPort + + state, err := pm.Status() + if err != nil { + result["phase"] = "ERROR" + result["error"] = err.Error() + return result + } + + result["phase"] = state.Phase + result["pid"] = state.PID + result["alive"] = state.Alive + result["port_ready"] = state.PortReady + result["frp_ready"] = state.FRPReady + if state.PortError != "" { + result["port_error"] = state.PortError + } + if state.Error != "" { + result["error"] = state.Error + } + if state.Version != "" { + result["version"] = state.Version + } + + return result +} diff --git a/main.go b/main.go index 66cd92c..98f637e 100644 --- a/main.go +++ b/main.go @@ -64,7 +64,7 @@ func main() { log.Printf("⚠️ 启动 frpc 失败: %v", err) } else { if status, err := pm.Status(); err == nil { - log.Printf("✅ frpc 状态: %s", status.State) + log.Printf("✅ frpc 状态: %s", status.Phase) if status.PID > 0 { log.Printf(" PID: %d, 端口: %d", status.PID, status.Port) } @@ -94,8 +94,8 @@ func startWatchdog(pm *process.ProcessManager) { continue } - if status.State != "running" { - log.Printf("⚠️ frpc 进程已停止 (状态: %s),自动重启...", status.State) + if status.Phase != "RUNNING" { + log.Printf("⚠️ frpc 进程已停止 (状态: %s),自动重启...", status.Phase) ctx := context.Background() if err := pm.Start(ctx); err != nil { log.Printf("❌ 自动重启 frpc 失败: %v", err)