Files
frpc-console/internal/process/manager.go
T

1365 lines
36 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// internal/process/manager.go
// 新增 PhaseReloading 状态 + FRPReady 绑定 PID
// 嗷呜! ^_^
// 作者留:现在2.7-Preview首战告捷!前端检测正常,后端匹配达成,该热加载的都能上了!
// 代价嘛……之前极致的 20M + 15M ,现在变成了 22M + 37M ……
// 多点就多点嘛……总比低内存消耗下黑盒状态强嘛……
package process
import (
"bufio"
"context"
"encoding/json"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
"os/exec"
"path/filepath"
"runtime"
"strconv"
"strings"
"sync"
"syscall"
"time"
)
// ================================================================
// 常量定义
// ================================================================
const (
LockFileName = ".frpc.lock"
PortCheckTimeout = 500 * time.Millisecond
StopMaxWaitTime = 5 * time.Second
LockAcquireTimeout = 30 * time.Second
LockRetryInterval = 100 * time.Millisecond
APITimeout = 2 * time.Second
StartupMaxAttempts = 50
StartupRetryDelay = 200 * time.Millisecond
StartupTimeout = 10 * time.Second
HealthCheckInterval = 15 * time.Second
)
// ================================================================
// 进程状态定义
// ================================================================
type ProcessPhase string
const (
PhaseUnknown ProcessPhase = "UNKNOWN"
PhaseStarting ProcessPhase = "STARTING"
PhaseRunning ProcessPhase = "RUNNING"
PhaseDegraded ProcessPhase = "DEGRADED"
PhaseFailed ProcessPhase = "FAILED"
PhaseStopped ProcessPhase = "STOPPED"
PhaseStopping ProcessPhase = "STOPPING"
PhaseConflict ProcessPhase = "CONFLICT"
PhaseReloading ProcessPhase = "RELOADING" // 新增:reload 中间态
)
// ================================================================
// 数据结构
// ================================================================
type PortCheckResult struct {
Ready bool
Err error
PID int
Process string
}
type FrpcInstance struct {
PID int
ParentPID int
ExecPath string
CmdLine string
Owned bool
}
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"`
PortPID int `json:"port_pid"`
PortOwner string `json:"port_owner"`
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"`
Conflicts []FrpcInstance `json:"conflicts,omitempty"`
}
type PortStatus struct {
Port int `json:"port"`
Occupied bool `json:"occupied"`
PID int `json:"pid"`
IsFRPC bool `json:"is_frpc"`
ProcessCmd string `json:"process_cmd,omitempty"`
}
// FRPCStatus 匹配 frpc admin API /api/status 的实际返回结构
type FRPCStatus struct {
TCP []struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
} `json:"tcp"`
UDP []struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
} `json:"udp"`
HTTP []struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
} `json:"http"`
HTTPS []struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
} `json:"https"`
STCP []struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
} `json:"stcp"`
XTCP []struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
} `json:"xtcp"`
SUDP []struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
} `json:"sudp"`
}
type ConflictInfo struct {
HasConflict bool
Count int
Owned []FrpcInstance
Unknown []FrpcInstance
}
// ================================================================
// 全局变量
// ================================================================
var (
globalManager *ProcessManager
globalMu sync.Mutex
)
// ================================================================
// ProcessManager 主结构
// ================================================================
type ProcessManager struct {
mu sync.Mutex
dataDir string
configPath string
frpcBinPath string
adminPort int
startTime time.Time
expectedStop bool
exitCode int
exitMu sync.Mutex
lastOutput string
lastError string
outputMu sync.Mutex
lockFile *os.File
locked bool
currentPhase ProcessPhase
statusMu sync.RWMutex
// 用于 admin API 检测的 HTTP 客户端
httpClient *http.Client
// 当前实例的 run_id(从 admin API 获取)
runID string
}
// ================================================================
// 构造函数
// ================================================================
func NewManager(dataDir, configPath, frpcBinPath string) *ProcessManager {
pm := &ProcessManager{
dataDir: dataDir,
configPath: configPath,
frpcBinPath: frpcBinPath,
adminPort: 0,
currentPhase: PhaseUnknown,
httpClient: &http.Client{
Timeout: APITimeout,
},
}
return pm
}
func SetGlobalManager(pm *ProcessManager) {
globalMu.Lock()
defer globalMu.Unlock()
globalManager = pm
}
func GetGlobalManager() *ProcessManager {
globalMu.Lock()
defer globalMu.Unlock()
return globalManager
}
func (pm *ProcessManager) AdminPort() int {
return pm.adminPort
}
func (pm *ProcessManager) CurrentPhase() ProcessPhase {
pm.statusMu.RLock()
defer pm.statusMu.RUnlock()
return pm.currentPhase
}
func (pm *ProcessManager) CurrentPID() int {
pm.statusMu.RLock()
defer pm.statusMu.RUnlock()
return pm.readPIDFile()
}
func (pm *ProcessManager) setPhase(phase ProcessPhase) {
pm.statusMu.Lock()
defer pm.statusMu.Unlock()
if pm.currentPhase != phase {
log.Printf("[STATUS] %s → %s", pm.currentPhase, phase)
pm.currentPhase = phase
}
}
// ================================================================
// 配置读取
// ================================================================
func (pm *ProcessManager) LoadConfig() error {
content, err := os.ReadFile(pm.configPath)
if err != nil {
return fmt.Errorf("读取配置文件失败: %w", err)
}
log.Printf("[DEBUG] LoadConfig 读取到文件,长度: %d 字节", len(content))
preview := string(content)
if len(preview) > 600 {
preview = preview[:600] + "\n... (截断)"
}
log.Printf("[DEBUG] 文件内容预览:\n%s", preview)
if port := extractIntValue(string(content), "admin_port"); port > 0 {
pm.adminPort = port
log.Printf("[DEBUG] ✅ 从 admin_port 解析到端口: %d", port)
return nil
}
log.Printf("[DEBUG] admin_port 未找到,尝试解析 webServer")
if port := extractIntValueFromSection(string(content), "webServer", "port"); port > 0 {
pm.adminPort = port
log.Printf("[DEBUG] ✅ 从 webServer.port 解析到端口: %d", port)
return nil
}
log.Printf("[DEBUG] webServer.port 未找到")
if port := extractPortFromAddrSection(string(content), "webServer", "addr"); port > 0 {
pm.adminPort = port
log.Printf("[DEBUG] ✅ 从 webServer.addr 解析到端口: %d", port)
return nil
}
log.Printf("[DEBUG] webServer.addr 未找到或解析失败")
return fmt.Errorf("未找到 admin_port 或 webServer.port/addr 配置")
}
func extractIntValue(content, key string) int {
lines := strings.Split(content, "\n")
for _, line := range lines {
trimmed := strings.TrimSpace(line)
if strings.HasPrefix(trimmed, key) {
parts := strings.SplitN(trimmed, "=", 2)
if len(parts) == 2 {
val := strings.TrimSpace(parts[1])
val = strings.Trim(val, `"`)
if port, err := strconv.Atoi(val); err == nil && port > 0 {
return port
}
}
}
}
return 0
}
func extractIntValueFromSection(content, section, key string) int {
lines := strings.Split(content, "\n")
inSection := false
for _, line := range lines {
trimmed := strings.TrimSpace(line)
if strings.HasPrefix(trimmed, "[") && strings.HasSuffix(trimmed, "]") {
sectionName := strings.TrimSpace(strings.Trim(trimmed, "[]"))
inSection = strings.EqualFold(sectionName, section)
continue
}
if inSection && strings.HasPrefix(trimmed, key) {
parts := strings.SplitN(trimmed, "=", 2)
if len(parts) == 2 {
val := strings.TrimSpace(parts[1])
val = strings.Trim(val, `"`)
if port, err := strconv.Atoi(val); err == nil && port > 0 {
return port
}
}
}
}
return 0
}
func extractPortFromAddrSection(content, section, key string) int {
lines := strings.Split(content, "\n")
inSection := false
for _, line := range lines {
trimmed := strings.TrimSpace(line)
if strings.HasPrefix(trimmed, "[") && strings.HasSuffix(trimmed, "]") {
sectionName := strings.TrimSpace(strings.Trim(trimmed, "[]"))
inSection = strings.EqualFold(sectionName, section)
continue
}
if inSection && strings.HasPrefix(trimmed, key) {
parts := strings.SplitN(trimmed, "=", 2)
if len(parts) == 2 {
val := strings.TrimSpace(parts[1])
val = strings.Trim(val, `"`)
if idx := strings.LastIndex(val, ":"); idx != -1 {
portStr := val[idx+1:]
if port, err := strconv.Atoi(portStr); err == nil && port > 0 {
return port
}
}
}
}
}
return 0
}
// ================================================================
// 端口检测(含归属信息)
// ================================================================
func (pm *ProcessManager) CheckPort() PortCheckResult {
if pm.adminPort <= 0 {
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 PortCheckResult{Ready: false, Err: err}
}
conn.Close()
// 获取端口占用者信息
pid, process := pm.getPortOwner(pm.adminPort)
// 即使 pid=0,也返回 Ready=true,让调用方决定如何处理
return PortCheckResult{
Ready: true,
PID: pid,
Process: process,
}
}
func (pm *ProcessManager) getPortOwner(port int) (int, string) {
// 方法1: ss
if pid, name := pm.getPortOwnerBySS(port); pid > 0 {
return pid, name
}
// 方法2: netstat
if pid, name := pm.getPortOwnerByNetstat(port); pid > 0 {
return pid, name
}
return 0, ""
}
func (pm *ProcessManager) getPortOwnerBySS(port int) (int, string) {
// ss -lntp | grep ':7400 ' | grep -oP 'pid=\K[0-9]+' | head -1
cmd := exec.Command("sh", "-c", fmt.Sprintf("ss -lntp | grep ':%d ' | grep -oP 'pid=\\K[0-9]+' | head -1", port))
out, err := cmd.Output()
if err != nil {
return 0, ""
}
pidStr := strings.TrimSpace(string(out))
if pidStr == "" {
return 0, ""
}
pid, err := strconv.Atoi(pidStr)
if err != nil || pid <= 0 {
return 0, ""
}
return pid, ""
}
func (pm *ProcessManager) getPortOwnerByNetstat(port int) (int, string) {
// netstat -tlnp | grep ':7400 ' | awk '{print $7}' | cut -d'/' -f1 | head -1
cmd := exec.Command("sh", "-c", fmt.Sprintf("netstat -tlnp 2>/dev/null | grep ':%d ' | awk '{print $7}' | cut -d'/' -f1 | head -1", port))
out, err := cmd.Output()
if err != nil {
return 0, ""
}
pidStr := strings.TrimSpace(string(out))
if pidStr == "" {
return 0, ""
}
pid, err := strconv.Atoi(pidStr)
if err != nil || pid <= 0 {
return 0, ""
}
return pid, ""
}
// ================================================================
// 实例检测
// ================================================================
func (pm *ProcessManager) DetectFrpcInstances() []FrpcInstance {
var instances []FrpcInstance
if runtime.GOOS == "linux" {
cmd := exec.Command("pgrep", "-f", "frpc")
out, err := cmd.Output()
if err != nil {
return instances
}
pids := strings.Split(strings.TrimSpace(string(out)), "\n")
for _, pidStr := range pids {
pidStr = strings.TrimSpace(pidStr)
if pidStr == "" {
continue
}
pid, err := strconv.Atoi(pidStr)
if err != nil || pid <= 0 {
continue
}
if pid == 1 || pid == os.Getpid() {
continue
}
if !pm.isProcessAlive(pid) {
continue
}
cmdline, _ := os.ReadFile(fmt.Sprintf("/proc/%d/cmdline", pid))
cmdLineStr := strings.ReplaceAll(string(cmdline), "\x00", " ")
exePath, _ := os.Readlink(fmt.Sprintf("/proc/%d/exe", pid))
stat, _ := os.ReadFile(fmt.Sprintf("/proc/%d/stat", pid))
var parentPID int
if len(stat) > 0 {
parts := strings.Fields(string(stat))
if len(parts) > 3 {
parentPID, _ = strconv.Atoi(parts[3])
}
}
inst := FrpcInstance{
PID: pid,
ParentPID: parentPID,
ExecPath: exePath,
CmdLine: cmdLineStr,
Owned: false,
}
ownedPid := pm.readPIDFile()
if ownedPid == pid {
inst.Owned = true
}
instances = append(instances, inst)
}
}
return instances
}
func (pm *ProcessManager) filterOwned(instances []FrpcInstance) []FrpcInstance {
var result []FrpcInstance
for _, inst := range instances {
if inst.Owned {
result = append(result, inst)
}
}
return result
}
func (pm *ProcessManager) filterUnknown(instances []FrpcInstance) []FrpcInstance {
var result []FrpcInstance
for _, inst := range instances {
if !inst.Owned {
result = append(result, inst)
}
}
return result
}
func (pm *ProcessManager) DetectConflict() ConflictInfo {
instances := pm.DetectFrpcInstances()
owned := pm.filterOwned(instances)
unknown := pm.filterUnknown(instances)
return ConflictInfo{
HasConflict: len(unknown) > 0 || len(owned) > 1,
Count: len(instances),
Owned: owned,
Unknown: unknown,
}
}
// ================================================================
// 进程存活检测
// ================================================================
func (pm *ProcessManager) isProcessAlive(pid int) bool {
if pid <= 0 {
return false
}
if runtime.GOOS == "windows" {
cmd := exec.Command("tasklist", "/FI", "PID eq", strconv.Itoa(pid))
output, err := cmd.CombinedOutput()
if err != nil {
return false
}
return strings.Contains(string(output), strconv.Itoa(pid))
}
process, err := os.FindProcess(pid)
if err != nil {
return false
}
return process.Signal(syscall.Signal(0)) == nil
}
// ================================================================
// FRPReady 检测(绑定 PID
// ================================================================
func (pm *ProcessManager) isFRPReady(pid int) bool {
if pid <= 0 || pm.adminPort <= 0 {
log.Printf("[DEBUG] FRPReady(pid=%d): pid 无效或 admin_port 未配置", pid)
return false
}
// stdout 快速通道
if strings.Contains(pm.getLastOutput(), "start proxy success") ||
strings.Contains(pm.getLastOutput(), "login to server success") {
log.Printf("[DEBUG] FRPReady(pid=%d): 检测到 stdout 关键字", pid)
return true
}
// 端口归属检测
portResult := pm.CheckPort()
if !portResult.Ready {
log.Printf("[DEBUG] FRPReady(pid=%d): 端口 %d 未就绪", pid, pm.adminPort)
return false
}
// 关键修复:端口被占用但无法识别归属 → 保守返回 false
if portResult.PID == 0 {
log.Printf("[DEBUG] FRPReady(pid=%d): 端口 %d 被占用但无法识别归属进程,保守返回 false", pid, pm.adminPort)
return false
}
if portResult.PID != pid {
log.Printf("[DEBUG] FRPReady(pid=%d): 端口 %d 被进程 %d 占用,与期望 PID %d 不一致",
pid, pm.adminPort, portResult.PID, pid)
return false
}
// admin API 检测
url := fmt.Sprintf("http://127.0.0.1:%d/api/status", pm.adminPort)
ctx, cancel := context.WithTimeout(context.Background(), APITimeout)
defer cancel()
req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
if err != nil {
log.Printf("[DEBUG] FRPReady(pid=%d): 创建请求失败: %v", pid, err)
return false
}
resp, err := pm.httpClient.Do(req)
if err != nil {
log.Printf("[DEBUG] FRPReady(pid=%d): admin API 请求失败: %v", pid, err)
return false
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
log.Printf("[DEBUG] FRPReady(pid=%d): admin API 返回状态码 %d", pid, resp.StatusCode)
return false
}
body, err := io.ReadAll(resp.Body)
if err != nil {
log.Printf("[DEBUG] FRPReady(pid=%d): 读取响应失败: %v", pid, err)
return false
}
var status FRPCStatus
if err := json.Unmarshal(body, &status); err != nil {
log.Printf("[DEBUG] FRPReady(pid=%d): 解析 JSON 失败: %v", pid, err)
return false
}
// 检查代理状态
proxyLists := [][]struct {
Name string `json:"name"`
Type string `json:"type"`
Status string `json:"status"`
Err string `json:"err"`
LocalAddr string `json:"local_addr"`
RemoteAddr string `json:"remote_addr"`
}{
status.TCP, status.UDP, status.HTTP, status.HTTPS,
status.STCP, status.XTCP, status.SUDP,
}
for _, proxies := range proxyLists {
for _, p := range proxies {
if p.Status == "running" {
log.Printf("[DEBUG] FRPReady(pid=%d): 代理 %s 状态为 running", pid, p.Name)
return true
}
}
}
log.Printf("[DEBUG] FRPReady(pid=%d): 没有代理处于 running 状态", pid)
return false
}
// ================================================================
// 进程清理
// ================================================================
func (pm *ProcessManager) killProcess(pid int) error {
if pid <= 0 {
return nil
}
if runtime.GOOS == "windows" {
cmd := exec.Command("taskkill", "/F", "/PID", strconv.Itoa(pid))
return cmd.Run()
}
process, err := os.FindProcess(pid)
if err != nil {
return err
}
return process.Kill()
}
func (pm *ProcessManager) cleanupOrphans() {
result := pm.CheckPort()
if !result.Ready {
return
}
pid := pm.readPIDFile()
if pid > 0 && pm.isProcessAlive(pid) {
return
}
log.Printf("[WARN] 检测到孤儿 frpc 进程 (端口 %d 被占用但无有效 PID),正在清理...", pm.adminPort)
if runtime.GOOS == "windows" {
exec.Command("taskkill", "/F", "/IM", "frpc.exe").Run()
} else {
exec.Command("pkill", "-f", "frpc").Run()
}
pm.deletePIDFile()
time.Sleep(1 * time.Second)
}
// ================================================================
// PID 文件操作
// ================================================================
func (pm *ProcessManager) pidFilePath() string {
return filepath.Join(pm.dataDir, "frpc.pid")
}
func (pm *ProcessManager) readPIDFile() int {
data, err := os.ReadFile(pm.pidFilePath())
if err != nil {
return 0
}
pid, err := strconv.Atoi(strings.TrimSpace(string(data)))
if err != nil || pid <= 0 {
return 0
}
return pid
}
func (pm *ProcessManager) writePIDFile(pid int) error {
return os.WriteFile(pm.pidFilePath(), []byte(strconv.Itoa(pid)), 0644)
}
func (pm *ProcessManager) deletePIDFile() error {
err := os.Remove(pm.pidFilePath())
if os.IsNotExist(err) {
return nil
}
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
}
// ================================================================
// 状态计算
// ================================================================
func (pm *ProcessManager) computeState(pid int, result PortCheckResult) ProcessState {
alive := pm.isProcessAlive(pid)
portReady := result.Ready
portPID := result.PID
portOwner := result.Process
portErr := result.Err
pm.exitMu.Lock()
exitCode := pm.exitCode
pm.exitMu.Unlock()
state := ProcessState{
PID: pid,
Port: pm.adminPort,
Alive: alive,
PortReady: portReady,
PortPID: portPID,
PortOwner: portOwner,
PortError: "",
FRPReady: false,
StartedAt: pm.startTime,
ExpectedStop: pm.expectedStop,
LastOutput: pm.getLastOutput(),
LastError: pm.getLastError(),
ExitCode: exitCode,
}
if portErr != nil {
state.PortError = portErr.Error()
}
// ---- 状态判定逻辑 ----
// 1. 主动停止中
if pm.expectedStop && alive {
state.Phase = PhaseStopping
return state
}
// 2. 进程不存在
if !alive {
log.Printf("[DEBUG] computeState: alive=false, pid=%d, expectedStop=%v, exitCode=%d", pid, pm.expectedStop, exitCode)
if pm.expectedStop {
state.Phase = PhaseStopped
return state
}
if exitCode != 0 && exitCode != -1 {
state.Error = fmt.Sprintf("进程异常退出 (exit code: %d)", exitCode)
if state.LastError != "" {
state.Error += ": " + state.LastError
}
} else {
state.Error = "进程意外退出"
if state.LastError != "" {
state.Error += ": " + state.LastError
}
}
state.Phase = PhaseFailed
return state
}
// 3. 启动超时(只对 STARTING 状态生效)
if pm.currentPhase == PhaseStarting && time.Since(pm.startTime) > StartupTimeout {
log.Printf("[DEBUG] computeState: 启动超时, startTime=%v, elapsed=%v",
pm.startTime, time.Since(pm.startTime))
state.Error = fmt.Sprintf("启动超时 (超过 %v)", StartupTimeout)
if portErr != nil {
state.Error += ": " + portErr.Error()
}
state.Phase = PhaseFailed
return state
}
// 4. 端口检测结果
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
}
// 5. 端口就绪,检查归属
if portPID > 0 && pid > 0 && portPID != pid {
log.Printf("[DEBUG] computeState: 端口冲突, portPID=%d, pid=%d", portPID, pid)
// 端口被其他进程占用 → CONFLICT
state.Phase = PhaseConflict
state.Error = fmt.Sprintf("端口 %d 被进程 %d (%s) 占用,与 PID 文件 %d 不一致",
pm.adminPort, portPID, portOwner, pid)
state.FRPReady = false
return state
}
// 6. 端口被自己的进程占用,检查 FRPReady
state.FRPReady = pm.isFRPReady(pid)
if state.FRPReady {
state.Phase = PhaseRunning
return state
}
// 7. 端口就绪但 FRP 未就绪
state.Phase = PhaseDegraded
state.Error = "端口已监听,但 admin API 未就绪或代理未运行"
return state
}
// ================================================================
// 状态查询
// ================================================================
func (pm *ProcessManager) Status() (*ProcessState, error) {
pid := pm.readPIDFile()
result := pm.CheckPort()
state := pm.computeState(pid, result)
pm.setPhase(state.Phase)
return &state, nil
}
// ================================================================
// Reload 操作(P0 核心)
// ================================================================
// ReloadConfig 执行 frpc 配置热加载,带状态管理
func (pm *ProcessManager) ReloadConfig(ctx context.Context) error {
if err := pm.Lock(); err != nil {
return fmt.Errorf("获取锁失败: %w", err)
}
defer pm.Unlock()
// 获取当前状态
pid := pm.readPIDFile()
if pid <= 0 || !pm.isProcessAlive(pid) {
// 进程不存在,直接启动
log.Printf("[INFO] Reload: frpc 未运行,执行启动")
return pm.startLocked(ctx)
}
// 记录 reload 前的 PID 和 run_id
oldPID := pid
log.Printf("[INFO] Reload: 开始热加载 (当前 PID: %d)", oldPID)
// 进入 RELOADING 状态
pm.setPhase(PhaseReloading)
// 执行 frpc reload 命令
frpcPath, err := pm.getFrpcPath()
if err != nil {
pm.setPhase(PhaseDegraded)
return fmt.Errorf("获取 frpc 路径失败: %w", err)
}
cmd := exec.Command(frpcPath, "reload", "-c", pm.configPath)
output, err := cmd.CombinedOutput()
// 检查 reload 执行结果
if err != nil {
// reload 命令失败,检查进程是否还在
if !pm.isProcessAlive(oldPID) {
// 进程已退出,reload 失败且进程丢失
log.Printf("[WARN] Reload: frpc 进程在 reload 期间退出 (PID: %d)", oldPID)
pm.setPhase(PhaseFailed)
pm.deletePIDFile()
return fmt.Errorf("reload 失败,frpc 进程已退出: %w", err)
}
// 进程还在,但 reload 命令失败,可能是配置问题
log.Printf("[WARN] Reload: 命令失败但进程仍在运行 (PID: %d), 输出: %s", oldPID, string(output))
pm.setPhase(PhaseDegraded)
return fmt.Errorf("reload 命令执行失败: %w", err)
}
log.Printf("[INFO] Reload: 命令执行成功,输出: %s", string(output))
// 等待新进程就绪
time.Sleep(1 * time.Second)
// 获取新进程的 PID
newPID := pm.readPIDFile()
if newPID <= 0 || newPID == oldPID {
// PID 没变化,可能是 reload 没有触发进程切换
log.Printf("[INFO] Reload: PID 未变化 (PID: %d),验证服务状态...", oldPID)
if pm.isFRPReady(oldPID) {
pm.setPhase(PhaseRunning)
log.Printf("[INFO] Reload: 服务仍健康,保持运行 (PID: %d)", oldPID)
return nil
}
pm.setPhase(PhaseDegraded)
return fmt.Errorf("reload 后服务未就绪 (PID: %d)", oldPID)
}
// PID 已变化,验证新进程
log.Printf("[INFO] Reload: PID 从 %d 变为 %d", oldPID, newPID)
// 等待新进程的 FRPReady
for attempt := 0; attempt < 20; attempt++ {
if pm.isFRPReady(newPID) {
pm.setPhase(PhaseRunning)
log.Printf("[INFO] Reload: 成功切换到新进程 (PID: %d)", newPID)
return nil
}
time.Sleep(200 * time.Millisecond)
}
// 新进程未就绪,回退状态
pm.setPhase(PhaseDegraded)
return fmt.Errorf("reload 后新进程未就绪 (PID: %d)", newPID)
}
// getFrpcPath 获取 frpc 二进制路径
func (pm *ProcessManager) getFrpcPath() (string, error) {
// 如果 frpcBinPath 有效,直接返回
if pm.frpcBinPath != "" {
if _, err := os.Stat(pm.frpcBinPath); err == nil {
return pm.frpcBinPath, nil
}
}
// 否则使用 frp 模块的 GetFrpcPath
// 避免循环引用,从外部传入
return pm.frpcBinPath, nil
}
// ================================================================
// 操作执行
// ================================================================
// 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.expectedStop = false
conflict := pm.DetectConflict()
if conflict.HasConflict {
log.Printf("[WARN] 检测到 %d 个 frpc 实例 (Owned: %d, Unknown: %d),进入 CONFLICT 状态",
conflict.Count, len(conflict.Owned), len(conflict.Unknown))
for _, inst := range conflict.Unknown {
log.Printf("[INFO] 清理 Unknown frpc 实例 (PID: %d)", inst.PID)
pm.killProcess(inst.PID)
}
if len(conflict.Owned) == 1 {
pid := conflict.Owned[0].PID
log.Printf("[INFO] 保留 Owned frpc 实例 (PID: %d),清理完成", pid)
pm.writePIDFile(pid)
pm.setPhase(PhaseRunning)
return nil
}
if len(conflict.Owned) > 1 {
log.Printf("[WARN] 多个 Owned 实例冲突,全部清理")
for _, inst := range conflict.Owned {
pm.killProcess(inst.PID)
}
}
pm.deletePIDFile()
}
pm.cleanupOrphans()
return pm.startLocked(ctx)
}
func (pm *ProcessManager) startLocked(ctx context.Context) error {
result := pm.CheckPort()
if result.Ready {
pid := pm.readPIDFile()
if pid > 0 && pm.isProcessAlive(pid) {
if result.PID == pid {
log.Printf("[INFO] frpc 已在运行 (PID: %d)", pid)
pm.setPhase(PhaseRunning)
return nil
}
log.Printf("[WARN] 端口被进程 %d 占用,但 PID 文件指向 %d,可能存在冲突", result.PID, pid)
}
pm.cleanupOrphans()
}
pm.startTime = time.Now()
pm.expectedStop = false
pm.runID = ""
pm.exitMu.Lock()
pm.exitCode = -1
pm.exitMu.Unlock()
cmd := exec.CommandContext(ctx, pm.frpcBinPath, "-c", pm.configPath)
setProcessAttributes(cmd)
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)
}
if err := pm.writePIDFile(cmd.Process.Pid); err != nil {
log.Printf("[WARN] 写入 PID 文件失败: %v", err)
}
log.Printf("[DEBUG] frpc 进程已启动,PID: %d", cmd.Process.Pid)
pm.setPhase(PhaseStarting)
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.currentPhase != PhaseReloading {
pm.setPhase(PhaseFailed)
}
}
if !pm.expectedStop {
pm.deletePIDFile()
}
}()
for attempt := 0; attempt < StartupMaxAttempts; attempt++ {
r := pm.CheckPort()
if r.Ready {
if r.PID > 0 && r.PID != cmd.Process.Pid {
log.Printf("[WARN] 端口被进程 %d 占用,与当前进程 %d 不一致,可能存在冲突", r.PID, cmd.Process.Pid)
time.Sleep(StartupRetryDelay)
continue
}
if pm.isFRPReady(cmd.Process.Pid) {
log.Printf("[INFO] frpc 启动成功 (PID: %d, 端口: %d),耗时 %dms",
cmd.Process.Pid, pm.adminPort, attempt*int(StartupRetryDelay/time.Millisecond))
pm.setPhase(PhaseRunning)
return nil
}
log.Printf("[DEBUG] 端口已就绪,等待 frpc 初始化...")
}
time.Sleep(StartupRetryDelay)
}
if pm.isProcessAlive(cmd.Process.Pid) {
stderr := pm.getLastError()
log.Printf("[WARN] frpc 启动超时 (PID: %d),当前 stderr: %s", cmd.Process.Pid, stderr)
if pm.isFRPReady(cmd.Process.Pid) {
log.Printf("[INFO] frpc 实际已就绪 (admin API 可访问),超时误判,修正状态")
pm.setPhase(PhaseRunning)
return nil
}
pm.killProcess(cmd.Process.Pid)
pm.deletePIDFile()
pm.setPhase(PhaseFailed)
return fmt.Errorf("frpc 启动超时: 进程存在但端口未就绪")
}
pm.deletePIDFile()
pm.setPhase(PhaseFailed)
return fmt.Errorf("frpc 启动失败: 进程已退出")
}
func (pm *ProcessManager) Stop(ctx context.Context) error {
if err := pm.Lock(); err != nil {
return fmt.Errorf("获取锁失败: %w", err)
}
defer pm.Unlock()
pm.expectedStop = true
pm.setPhase(PhaseStopping)
return pm.stopLocked(ctx)
}
func (pm *ProcessManager) stopLocked(_ context.Context) error {
pid := pm.readPIDFile()
if pid <= 0 {
if runtime.GOOS == "windows" {
exec.Command("taskkill", "/F", "/IM", "frpc.exe").Run()
} else {
exec.Command("pkill", "-f", "frpc").Run()
}
pm.deletePIDFile()
pm.setPhase(PhaseStopped)
return nil
}
if !pm.isProcessAlive(pid) {
pm.deletePIDFile()
pm.setPhase(PhaseStopped)
return nil
}
process, err := os.FindProcess(pid)
if err != nil {
pm.deletePIDFile()
pm.setPhase(PhaseStopped)
return nil
}
if err := process.Signal(syscall.SIGTERM); err != nil {
pm.deletePIDFile()
pm.setPhase(PhaseStopped)
return nil
}
start := time.Now()
for time.Since(start) < StopMaxWaitTime {
if !pm.isProcessAlive(pid) {
pm.deletePIDFile()
pm.setPhase(PhaseStopped)
return nil
}
time.Sleep(200 * time.Millisecond)
}
if runtime.GOOS == "windows" {
exec.Command("taskkill", "/F", "/PID", strconv.Itoa(pid)).Run()
} else {
process.Kill()
}
time.Sleep(500 * time.Millisecond)
if pm.isProcessAlive(pid) {
pm.setPhase(PhaseFailed)
return fmt.Errorf("强制停止失败: 进程仍存活 (PID: %d)", pid)
}
pm.deletePIDFile()
pm.setPhase(PhaseStopped)
return nil
}
func (pm *ProcessManager) Restart(ctx context.Context) error {
if err := pm.Lock(); err != nil {
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) StartHealthMonitor(ctx context.Context) {
ticker := time.NewTicker(HealthCheckInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
pm.runHealthCheck()
}
}
}
func (pm *ProcessManager) runHealthCheck() {
// 如果正在 reloading,跳过健康检查
if pm.currentPhase == PhaseReloading {
log.Printf("[DEBUG] 健康检查跳过: 正在 RELOADING")
return
}
state, err := pm.Status()
if err != nil {
log.Printf("[WARN] 健康检查失败: %v", err)
return
}
switch state.Phase {
case PhaseConflict:
log.Printf("[WARN] 检测到冲突,尝试自动恢复...")
for _, inst := range state.Conflicts {
if !inst.Owned {
pm.killProcess(inst.PID)
}
}
pm.setPhase(PhaseStarting)
case PhaseDegraded:
log.Printf("[WARN] frpc 处于降级状态,尝试恢复...")
pm.Restart(context.Background())
case PhaseFailed:
if !state.ExpectedStop {
log.Printf("[WARN] frpc 已失败,自动重启...")
pm.Start(context.Background())
}
}
}
// ================================================================
// 健康检查(供 API 调用)
// ================================================================
func (pm *ProcessManager) HealthCheck() map[string]interface{} {
result := make(map[string]interface{})
result["admin_port"] = pm.adminPort
result["phase"] = pm.CurrentPhase()
state, err := pm.Status()
if err != nil {
result["error"] = err.Error()
return result
}
result["pid"] = state.PID
result["port"] = state.Port
result["alive"] = state.Alive
result["port_ready"] = state.PortReady
result["port_pid"] = state.PortPID
result["port_owner"] = state.PortOwner
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
}
if len(state.Conflicts) > 0 {
result["conflicts"] = state.Conflicts
}
return result
}