mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-05 23:26:38 +08:00
[优化] 添加 Relay frps 连接和代理计数字段,重构相关逻辑以支持监控和可观测性
This commit is contained in:
@@ -1,5 +1,9 @@
|
||||
ARG VERSION=dev
|
||||
|
||||
FROM golang:1.25-alpine AS builder
|
||||
|
||||
ARG VERSION
|
||||
|
||||
WORKDIR /build
|
||||
|
||||
COPY openflare_relay/go.mod openflare_relay/go.sum ./
|
||||
@@ -7,7 +11,7 @@ COPY openflare_server /openflare_server
|
||||
COPY openflare_relay /openflare_relay
|
||||
|
||||
WORKDIR /openflare_relay
|
||||
RUN CGO_ENABLED=0 GOOS=linux go build -o openflare-relay ./cmd/relay
|
||||
RUN CGO_ENABLED=0 GOOS=linux go build -trimpath -ldflags "-s -w -X 'openflare-relay/internal/config.Version=$VERSION'" -o openflare-relay ./cmd/relay
|
||||
|
||||
# Final runtime image
|
||||
FROM fatedier/frps:v0.69.0
|
||||
|
||||
@@ -57,7 +57,7 @@ func main() {
|
||||
FrpsManager: frpsManager,
|
||||
HttpClient: httpClient,
|
||||
WebSocketService: wsClient,
|
||||
HeartbeatService: heartbeat.New(httpClient, frpsManager, cfg),
|
||||
HeartbeatService: heartbeat.New(httpClient, frpsManager, cfg, stateStore),
|
||||
}
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
package config
|
||||
|
||||
var Version = "dev"
|
||||
@@ -24,6 +24,17 @@ type Manager struct {
|
||||
activeConfig *service.RelayConfig
|
||||
cmd *exec.Cmd
|
||||
status string
|
||||
lastError string
|
||||
generation uint64
|
||||
stopping bool
|
||||
}
|
||||
|
||||
type RuntimeStatus struct {
|
||||
Status string
|
||||
LastError string
|
||||
Connections int
|
||||
ProxyCount int
|
||||
ProcessAlive bool
|
||||
}
|
||||
|
||||
func NewManager(frpsPath string, dataDir string) *Manager {
|
||||
@@ -54,6 +65,18 @@ func (m *Manager) GetStatus() string {
|
||||
return m.status
|
||||
}
|
||||
|
||||
func (m *Manager) GetRuntimeStatus() RuntimeStatus {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
return RuntimeStatus{
|
||||
Status: m.status,
|
||||
LastError: m.lastError,
|
||||
Connections: 0,
|
||||
ProxyCount: 0,
|
||||
ProcessAlive: m.cmd != nil && m.cmd.Process != nil,
|
||||
}
|
||||
}
|
||||
|
||||
func (m *Manager) UpdateConfig(cfg *service.RelayConfig) {
|
||||
if cfg == nil {
|
||||
return
|
||||
@@ -66,23 +89,35 @@ func (m *Manager) UpdateConfig(cfg *service.RelayConfig) {
|
||||
m.activeConfig.BindPort == cfg.BindPort &&
|
||||
m.activeConfig.VhostHTTPPort == cfg.VhostHTTPPort &&
|
||||
m.activeConfig.AuthToken == cfg.AuthToken {
|
||||
return // No change
|
||||
if m.cmd == nil && !m.stopping {
|
||||
slog.Warn("frps config unchanged but process is not running, restarting")
|
||||
if err := m.restartProcess(); err != nil {
|
||||
m.status = "unhealthy"
|
||||
m.lastError = err.Error()
|
||||
slog.Error("failed to restart frps with unchanged config", "error", err)
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
m.activeConfig = cfg
|
||||
m.stopping = false
|
||||
slog.Info("relay config updated, reloading frps")
|
||||
|
||||
if err := m.renderConfig(cfg); err != nil {
|
||||
slog.Error("failed to render frps config", "error", err)
|
||||
m.status = "unhealthy"
|
||||
m.lastError = err.Error()
|
||||
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 = ""
|
||||
}
|
||||
}
|
||||
|
||||
@@ -105,13 +140,17 @@ func (m *Manager) renderConfig(cfg *service.RelayConfig) error {
|
||||
}
|
||||
|
||||
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.Wait()
|
||||
m.cmd = nil
|
||||
}
|
||||
return m.startProcessLocked(generation)
|
||||
}
|
||||
|
||||
func (m *Manager) startProcessLocked(generation uint64) error {
|
||||
cmd := exec.Command(m.frpsPath, "-c", m.configPath)
|
||||
cmd.Stdout = os.Stdout
|
||||
cmd.Stderr = os.Stderr
|
||||
@@ -121,8 +160,9 @@ func (m *Manager) restartProcess() error {
|
||||
}
|
||||
|
||||
m.cmd = cmd
|
||||
m.status = "healthy"
|
||||
m.lastError = ""
|
||||
|
||||
// Start a goroutine to monitor process exit
|
||||
go func(c *exec.Cmd) {
|
||||
err := c.Wait()
|
||||
slog.Warn("frps process exited", "error", err)
|
||||
@@ -130,8 +170,29 @@ func (m *Manager) restartProcess() error {
|
||||
if m.cmd == c {
|
||||
m.cmd = nil
|
||||
m.status = "unhealthy"
|
||||
if err != nil {
|
||||
m.lastError = err.Error()
|
||||
} else {
|
||||
m.lastError = "frps process exited"
|
||||
}
|
||||
}
|
||||
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
|
||||
@@ -140,9 +201,11 @@ func (m *Manager) restartProcess() error {
|
||||
func (m *Manager) Stop() {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.stopping = true
|
||||
m.generation++
|
||||
if m.cmd != nil && m.cmd.Process != nil {
|
||||
_ = m.cmd.Process.Kill()
|
||||
_ = m.cmd.Wait()
|
||||
m.cmd = nil
|
||||
}
|
||||
m.status = "unhealthy"
|
||||
}
|
||||
|
||||
@@ -8,6 +8,8 @@ import (
|
||||
"openflare-relay/internal/config"
|
||||
"openflare-relay/internal/frps"
|
||||
"openflare-relay/internal/httpclient"
|
||||
"openflare-relay/internal/observability"
|
||||
"openflare-relay/internal/state"
|
||||
"openflare/service"
|
||||
)
|
||||
|
||||
@@ -15,13 +17,15 @@ type Service struct {
|
||||
client *httpclient.Client
|
||||
frpsManager *frps.Manager
|
||||
config *config.Config
|
||||
stateStore *state.Store
|
||||
}
|
||||
|
||||
func New(client *httpclient.Client, manager *frps.Manager, cfg *config.Config) *Service {
|
||||
func New(client *httpclient.Client, manager *frps.Manager, cfg *config.Config, stateStore *state.Store) *Service {
|
||||
return &Service{
|
||||
client: client,
|
||||
frpsManager: manager,
|
||||
config: cfg,
|
||||
stateStore: stateStore,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,12 +49,18 @@ func (s *Service) Run(ctx context.Context) {
|
||||
func (s *Service) doHeartbeat(ctx context.Context) {
|
||||
slog.Debug("sending heartbeat")
|
||||
|
||||
runtimeStatus := s.frpsManager.GetRuntimeStatus()
|
||||
payload := service.RelayHeartbeatPayload{
|
||||
RelayVersion: "0.1.0", // TODO dynamically inject build version
|
||||
RelayVersion: config.Version,
|
||||
FrpVersion: s.frpsManager.GetVersion(),
|
||||
RelayStatus: s.frpsManager.GetStatus(),
|
||||
FrpsConnCount: 0,
|
||||
FrpsProxyCount: 0,
|
||||
RelayStatus: runtimeStatus.Status,
|
||||
FrpsConnCount: runtimeStatus.Connections,
|
||||
FrpsProxyCount: runtimeStatus.ProxyCount,
|
||||
Name: s.config.NodeName,
|
||||
IP: s.config.NodeIP,
|
||||
Profile: observability.BuildProfile(s.config, s.stateStore),
|
||||
Snapshot: observability.BuildSnapshot(s.config, s.stateStore),
|
||||
HealthEvents: observability.BuildHealthEvents(runtimeStatus),
|
||||
}
|
||||
|
||||
resp, err := s.client.Heartbeat(ctx, payload)
|
||||
|
||||
@@ -0,0 +1,325 @@
|
||||
package observability
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"strconv"
|
||||
"strings"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"openflare-relay/internal/config"
|
||||
"openflare-relay/internal/frps"
|
||||
"openflare-relay/internal/state"
|
||||
"openflare/service"
|
||||
)
|
||||
|
||||
func BuildProfile(cfg *config.Config, stateStore *state.Store) *service.AgentNodeSystemProfile {
|
||||
profile := collectProfile(cfg)
|
||||
if profile == nil || stateStore == nil {
|
||||
return profile
|
||||
}
|
||||
fingerprint := fingerprintProfile(profile)
|
||||
snapshot, err := stateStore.Load()
|
||||
if err != nil {
|
||||
return profile
|
||||
}
|
||||
if snapshot.LastProfileFingerprint == fingerprint {
|
||||
return nil
|
||||
}
|
||||
snapshot.LastProfileFingerprint = fingerprint
|
||||
if err = stateStore.Save(snapshot); err != nil {
|
||||
return profile
|
||||
}
|
||||
return profile
|
||||
}
|
||||
|
||||
func BuildSnapshot(cfg *config.Config, stateStore *state.Store) *service.AgentNodeMetricSnapshot {
|
||||
now := time.Now().UTC()
|
||||
metric := &service.AgentNodeMetricSnapshot{CapturedAtUnix: now.Unix()}
|
||||
|
||||
metric.MemoryTotalBytes, metric.MemoryUsedBytes = readMemInfo()
|
||||
metric.StorageTotalBytes, metric.StorageUsedBytes = statFilesystem(cfg.DataDir)
|
||||
metric.NetworkRxBytes, metric.NetworkTxBytes = readLinuxNetworkTotals()
|
||||
metric.DiskReadBytes, metric.DiskWriteBytes = readLinuxDiskTotals()
|
||||
|
||||
if stateStore == nil {
|
||||
return metric
|
||||
}
|
||||
totalCPU, idleCPU := readLinuxCPUStat()
|
||||
snapshot, err := stateStore.Load()
|
||||
if err != nil {
|
||||
return metric
|
||||
}
|
||||
if snapshot.LastCPUStatTotal > 0 && totalCPU > snapshot.LastCPUStatTotal && idleCPU >= snapshot.LastCPUStatIdle {
|
||||
deltaTotal := totalCPU - snapshot.LastCPUStatTotal
|
||||
deltaIdle := idleCPU - snapshot.LastCPUStatIdle
|
||||
if deltaTotal > 0 && deltaIdle <= deltaTotal {
|
||||
metric.CPUUsagePercent = float64(deltaTotal-deltaIdle) / float64(deltaTotal) * 100
|
||||
}
|
||||
}
|
||||
snapshot.LastCPUStatTotal = totalCPU
|
||||
snapshot.LastCPUStatIdle = idleCPU
|
||||
snapshot.LastMetricAtUnix = now.Unix()
|
||||
_ = stateStore.Save(snapshot)
|
||||
return metric
|
||||
}
|
||||
|
||||
func BuildHealthEvents(status frps.RuntimeStatus) []service.AgentNodeHealthEvent {
|
||||
if strings.TrimSpace(status.Status) == "healthy" {
|
||||
return []service.AgentNodeHealthEvent{}
|
||||
}
|
||||
message := strings.TrimSpace(status.LastError)
|
||||
if message == "" {
|
||||
message = "frps runtime is not healthy"
|
||||
}
|
||||
return []service.AgentNodeHealthEvent{{
|
||||
EventType: "frps_unhealthy",
|
||||
Severity: "critical",
|
||||
Message: message,
|
||||
TriggeredAtUnix: time.Now().UTC().Unix(),
|
||||
}}
|
||||
}
|
||||
|
||||
func collectProfile(cfg *config.Config) *service.AgentNodeSystemProfile {
|
||||
hostname, _ := os.Hostname()
|
||||
osName, osVersion := readLinuxOSRelease()
|
||||
totalMemory, _ := readMemInfo()
|
||||
totalDisk, _ := statFilesystem(cfg.DataDir)
|
||||
return &service.AgentNodeSystemProfile{
|
||||
Hostname: strings.TrimSpace(hostname),
|
||||
OSName: osName,
|
||||
OSVersion: osVersion,
|
||||
KernelVersion: readFirstLine("/proc/sys/kernel/osrelease"),
|
||||
Architecture: runtime.GOARCH,
|
||||
CPUModel: readLinuxCPUModel(),
|
||||
CPUCores: runtime.NumCPU(),
|
||||
TotalMemoryBytes: totalMemory,
|
||||
TotalDiskBytes: totalDisk,
|
||||
UptimeSeconds: readLinuxUptimeSeconds(),
|
||||
ReportedAtUnix: time.Now().UTC().Unix(),
|
||||
}
|
||||
}
|
||||
|
||||
func fingerprintProfile(profile *service.AgentNodeSystemProfile) string {
|
||||
raw, err := json.Marshal(profile)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
sum := sha256.Sum256(raw)
|
||||
return hex.EncodeToString(sum[:])
|
||||
}
|
||||
|
||||
func readLinuxOSRelease() (string, string) {
|
||||
file, err := os.Open("/etc/os-release")
|
||||
if err != nil {
|
||||
return runtime.GOOS, ""
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
values := make(map[string]string)
|
||||
scanner := bufio.NewScanner(file)
|
||||
for scanner.Scan() {
|
||||
key, value, ok := strings.Cut(strings.TrimSpace(scanner.Text()), "=")
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
values[key] = strings.Trim(value, `"`)
|
||||
}
|
||||
if pretty := strings.TrimSpace(values["PRETTY_NAME"]); pretty != "" {
|
||||
return pretty, strings.TrimSpace(values["VERSION_ID"])
|
||||
}
|
||||
if name := strings.TrimSpace(values["NAME"]); name != "" {
|
||||
return name, strings.TrimSpace(values["VERSION_ID"])
|
||||
}
|
||||
return runtime.GOOS, ""
|
||||
}
|
||||
|
||||
func readLinuxCPUModel() string {
|
||||
file, err := os.Open("/proc/cpuinfo")
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
scanner := bufio.NewScanner(file)
|
||||
for scanner.Scan() {
|
||||
line := scanner.Text()
|
||||
if strings.HasPrefix(strings.ToLower(line), "model name") {
|
||||
_, value, ok := strings.Cut(line, ":")
|
||||
if ok {
|
||||
return strings.TrimSpace(value)
|
||||
}
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func readMemInfo() (int64, int64) {
|
||||
file, err := os.Open("/proc/meminfo")
|
||||
if err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
var totalKB, availableKB int64
|
||||
scanner := bufio.NewScanner(file)
|
||||
for scanner.Scan() {
|
||||
line := scanner.Text()
|
||||
if strings.HasPrefix(line, "MemTotal:") {
|
||||
totalKB = parseMemInfoValue(line)
|
||||
}
|
||||
if strings.HasPrefix(line, "MemAvailable:") {
|
||||
availableKB = parseMemInfoValue(line)
|
||||
}
|
||||
}
|
||||
total := totalKB * 1024
|
||||
used := total - availableKB*1024
|
||||
if used < 0 {
|
||||
used = 0
|
||||
}
|
||||
return total, used
|
||||
}
|
||||
|
||||
func parseMemInfoValue(line string) int64 {
|
||||
fields := strings.Fields(line)
|
||||
if len(fields) < 2 {
|
||||
return 0
|
||||
}
|
||||
value, err := strconv.ParseInt(fields[1], 10, 64)
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
return value
|
||||
}
|
||||
|
||||
func readLinuxUptimeSeconds() int64 {
|
||||
content, err := os.ReadFile("/proc/uptime")
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
fields := strings.Fields(string(content))
|
||||
if len(fields) == 0 {
|
||||
return 0
|
||||
}
|
||||
value, err := strconv.ParseFloat(fields[0], 64)
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
return int64(value)
|
||||
}
|
||||
|
||||
func readLinuxCPUStat() (uint64, uint64) {
|
||||
content, err := os.ReadFile("/proc/stat")
|
||||
if err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
for _, line := range strings.Split(string(content), "\n") {
|
||||
if !strings.HasPrefix(line, "cpu ") {
|
||||
continue
|
||||
}
|
||||
fields := strings.Fields(line)
|
||||
if len(fields) < 5 {
|
||||
return 0, 0
|
||||
}
|
||||
var total uint64
|
||||
for index := 1; index < len(fields); index++ {
|
||||
value, err := strconv.ParseUint(fields[index], 10, 64)
|
||||
if err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
total += value
|
||||
}
|
||||
idle, err := strconv.ParseUint(fields[4], 10, 64)
|
||||
if err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
return total, idle
|
||||
}
|
||||
return 0, 0
|
||||
}
|
||||
|
||||
func readLinuxNetworkTotals() (int64, int64) {
|
||||
file, err := os.Open("/proc/net/dev")
|
||||
if err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
var rx, tx int64
|
||||
scanner := bufio.NewScanner(file)
|
||||
for scanner.Scan() {
|
||||
name, data, ok := strings.Cut(strings.TrimSpace(scanner.Text()), ":")
|
||||
if !ok || strings.TrimSpace(name) == "lo" {
|
||||
continue
|
||||
}
|
||||
fields := strings.Fields(data)
|
||||
if len(fields) < 16 {
|
||||
continue
|
||||
}
|
||||
if value, err := strconv.ParseInt(fields[0], 10, 64); err == nil {
|
||||
rx += value
|
||||
}
|
||||
if value, err := strconv.ParseInt(fields[8], 10, 64); err == nil {
|
||||
tx += value
|
||||
}
|
||||
}
|
||||
return rx, tx
|
||||
}
|
||||
|
||||
func readLinuxDiskTotals() (int64, int64) {
|
||||
file, err := os.Open("/proc/diskstats")
|
||||
if err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
var readBytes, writeBytes int64
|
||||
scanner := bufio.NewScanner(file)
|
||||
for scanner.Scan() {
|
||||
fields := strings.Fields(scanner.Text())
|
||||
if len(fields) < 14 || shouldSkipDiskDevice(fields[2]) {
|
||||
continue
|
||||
}
|
||||
if value, err := strconv.ParseInt(fields[5], 10, 64); err == nil {
|
||||
readBytes += value * 512
|
||||
}
|
||||
if value, err := strconv.ParseInt(fields[9], 10, 64); err == nil {
|
||||
writeBytes += value * 512
|
||||
}
|
||||
}
|
||||
return readBytes, writeBytes
|
||||
}
|
||||
|
||||
func shouldSkipDiskDevice(device string) bool {
|
||||
return device == "" || strings.HasPrefix(device, "loop") || strings.HasPrefix(device, "ram") || strings.HasPrefix(device, "dm-")
|
||||
}
|
||||
|
||||
func statFilesystem(path string) (int64, int64) {
|
||||
if strings.TrimSpace(path) == "" {
|
||||
path = string(os.PathSeparator)
|
||||
}
|
||||
var stat syscall.Statfs_t
|
||||
if err := syscall.Statfs(filepath.Clean(path), &stat); err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
total := int64(stat.Blocks) * int64(stat.Bsize)
|
||||
used := total - int64(stat.Bavail)*int64(stat.Bsize)
|
||||
if used < 0 {
|
||||
used = 0
|
||||
}
|
||||
return total, used
|
||||
}
|
||||
|
||||
func readFirstLine(path string) string {
|
||||
content, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
return strings.TrimSpace(string(content))
|
||||
}
|
||||
@@ -13,7 +13,11 @@ type Store struct {
|
||||
}
|
||||
|
||||
type State struct {
|
||||
LastAuthToken string `json:"last_auth_token"`
|
||||
LastAuthToken string `json:"last_auth_token"`
|
||||
LastProfileFingerprint string `json:"last_profile_fingerprint"`
|
||||
LastCPUStatTotal uint64 `json:"last_cpu_stat_total"`
|
||||
LastCPUStatIdle uint64 `json:"last_cpu_stat_idle"`
|
||||
LastMetricAtUnix int64 `json:"last_metric_at_unix"`
|
||||
}
|
||||
|
||||
func NewStore(path string) *Store {
|
||||
|
||||
Reference in New Issue
Block a user