diff --git a/openflare_relay/Dockerfile b/openflare_relay/Dockerfile index 5ba2580c..08aad498 100644 --- a/openflare_relay/Dockerfile +++ b/openflare_relay/Dockerfile @@ -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 diff --git a/openflare_relay/cmd/relay/main.go b/openflare_relay/cmd/relay/main.go index 8f57abfd..2fcce189 100644 --- a/openflare_relay/cmd/relay/main.go +++ b/openflare_relay/cmd/relay/main.go @@ -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) diff --git a/openflare_relay/internal/config/version.go b/openflare_relay/internal/config/version.go new file mode 100644 index 00000000..f356490f --- /dev/null +++ b/openflare_relay/internal/config/version.go @@ -0,0 +1,3 @@ +package config + +var Version = "dev" diff --git a/openflare_relay/internal/frps/manager.go b/openflare_relay/internal/frps/manager.go index 8af5f44f..98348752 100644 --- a/openflare_relay/internal/frps/manager.go +++ b/openflare_relay/internal/frps/manager.go @@ -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" } diff --git a/openflare_relay/internal/heartbeat/service.go b/openflare_relay/internal/heartbeat/service.go index df05be8d..c8cd24a1 100644 --- a/openflare_relay/internal/heartbeat/service.go +++ b/openflare_relay/internal/heartbeat/service.go @@ -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) diff --git a/openflare_relay/internal/observability/collector.go b/openflare_relay/internal/observability/collector.go new file mode 100644 index 00000000..f0c0ad9d --- /dev/null +++ b/openflare_relay/internal/observability/collector.go @@ -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)) +} diff --git a/openflare_relay/internal/state/store.go b/openflare_relay/internal/state/store.go index 39c63699..84832803 100644 --- a/openflare_relay/internal/state/store.go +++ b/openflare_relay/internal/state/store.go @@ -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 { diff --git a/openflare_server/controller/relay.go b/openflare_server/controller/relay.go index ac622b19..a83effcd 100644 --- a/openflare_server/controller/relay.go +++ b/openflare_server/controller/relay.go @@ -26,6 +26,7 @@ func RelayHeartbeat(c *gin.Context) { if !bindJSON(c, &payload) { return } + payload.IP = service.ResolveReportedNodeIP(payload.IP, c.Request.RemoteAddr) authNode, ok := c.Get("relay_node") if !ok { respondUnauthorized(c, "无权进行此操作") diff --git a/openflare_server/model/migrate/v20.go b/openflare_server/model/migrate/v20.go new file mode 100644 index 00000000..6e312619 --- /dev/null +++ b/openflare_server/model/migrate/v20.go @@ -0,0 +1,52 @@ +// v20 records low-frequency Relay frps counters on nodes so the management UI +// can show whether the relay runtime is alive and reporting tunnel load. +package migrate + +import ( + "fmt" + + "gorm.io/gorm" +) + +type nodeV20 struct { + ID uint `gorm:"primaryKey"` + RelayFrpsConnections int `gorm:"column:relay_frps_connections"` + RelayFrpsProxyCount int `gorm:"column:relay_frps_proxy_count"` +} + +func (nodeV20) TableName() string { + return "nodes" +} + +func init() { + Register(V20()) +} + +func V20() Migration { + return Migration{ + FromVersion: 19, + ToVersion: 20, + Migrate: migrateV20, + Validate: validateV20, + } +} + +func migrateV20(ctx Context, db *gorm.DB, backend string) error { + if err := ctx.ApplyCurrentSchema(db, backend); err != nil { + return err + } + return validateV20(ctx, db, backend) +} + +func validateV20(ctx Context, db *gorm.DB, backend string) error { + if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 19); err != nil { + return err + } + if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_connections") { + return fmt.Errorf("column nodes.relay_frps_connections is missing") + } + if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_proxy_count") { + return fmt.Errorf("column nodes.relay_frps_proxy_count is missing") + } + return nil +} diff --git a/openflare_server/model/migrations.go b/openflare_server/model/migrations.go index 96352520..11645087 100644 --- a/openflare_server/model/migrations.go +++ b/openflare_server/model/migrations.go @@ -84,6 +84,8 @@ func (databaseSchemaMigrationContext) ValidateDatabaseSchemaVersion(db *gorm.DB, return validateDatabaseSchemaV18(db, backend) case 19: return validateDatabaseSchemaV19(db, backend) + case 20: + return validateDatabaseSchemaV20(db, backend) default: return fmt.Errorf("database schema validation for v%d is not defined", version) } @@ -1225,6 +1227,19 @@ func validateDatabaseSchemaV19(db *gorm.DB, backend string) error { return nil } +func validateDatabaseSchemaV20(db *gorm.DB, backend string) error { + if err := validateDatabaseSchemaV19(db, backend); err != nil { + return err + } + if !db.Migrator().HasColumn(&Node{}, "relay_frps_connections") { + return fmt.Errorf("column nodes.relay_frps_connections is missing") + } + if !db.Migrator().HasColumn(&Node{}, "relay_frps_proxy_count") { + return fmt.Errorf("column nodes.relay_frps_proxy_count is missing") + } + return nil +} + func databaseSchemaMigrations() []databaseSchemaMigration { ctx := databaseSchemaMigrationContext{} migrations := []databaseSchemaMigration{} diff --git a/openflare_server/model/node.go b/openflare_server/model/node.go index 2e8c98f1..e275a708 100644 --- a/openflare_server/model/node.go +++ b/openflare_server/model/node.go @@ -40,6 +40,8 @@ type Node struct { RelayStatus string `json:"relay_status" gorm:"size:16;not null;default:'unknown'"` RelayFrpVersion string `json:"relay_frp_version" gorm:"size:64"` RelayVersion string `json:"relay_version" gorm:"size:64"` + RelayFrpsConnections int `json:"relay_frps_connections" gorm:"not null;default:0"` + RelayFrpsProxyCount int `json:"relay_frps_proxy_count" gorm:"not null;default:0"` } func ListNodes() (nodes []*Node, err error) { diff --git a/openflare_server/service/agent.go b/openflare_server/service/agent.go index 39211f94..8b355a95 100644 --- a/openflare_server/service/agent.go +++ b/openflare_server/service/agent.go @@ -172,6 +172,8 @@ type NodeView struct { RelayStatus string `json:"relay_status"` RelayFrpVersion string `json:"relay_frp_version"` RelayVersion string `json:"relay_version"` + RelayFrpsConnections int `json:"relay_frps_connections"` + RelayFrpsProxyCount int `json:"relay_frps_proxy_count"` } func HeartbeatNode(node *model.Node, payload AgentNodePayload) (*HeartbeatResponse, error) { diff --git a/openflare_server/service/node.go b/openflare_server/service/node.go index cd2b8765..523c283b 100644 --- a/openflare_server/service/node.go +++ b/openflare_server/service/node.go @@ -355,6 +355,8 @@ func buildNodeView(node *model.Node) *NodeView { view.RelayStatus = node.RelayStatus view.RelayFrpVersion = node.RelayFrpVersion view.RelayVersion = node.RelayVersion + view.RelayFrpsConnections = node.RelayFrpsConnections + view.RelayFrpsProxyCount = node.RelayFrpsProxyCount return view } diff --git a/openflare_server/service/node_observability.go b/openflare_server/service/node_observability.go index 37e07773..11306591 100644 --- a/openflare_server/service/node_observability.go +++ b/openflare_server/service/node_observability.go @@ -27,6 +27,7 @@ type NodeObservabilityView struct { HealthEvents []*model.NodeHealthEvent `json:"health_events"` Analytics NodeObservabilityAnalytics `json:"analytics"` Trends NodeObservabilityTrends `json:"trends"` + RelayDashboard *RelayDashboardSnapshot `json:"relay_dashboard,omitempty"` } type NodeObservabilityAnalytics struct { @@ -47,6 +48,25 @@ type NodeHealthEventCleanupResult struct { DeletedCount int64 `json:"deleted_count"` } +type RelayDashboardSnapshot struct { + TotalProxies int `json:"total_proxies"` + OnlineProxies int `json:"online_proxies"` + OfflineProxies int `json:"offline_proxies"` + Proxies []RelayProxyStat `json:"proxies"` + TotalConnections int `json:"total_connections"` + ClientCounts int `json:"client_counts"` +} + +type RelayProxyStat struct { + Name string `json:"name"` + Type string `json:"type"` + Status string `json:"status"` + ClientVersion string `json:"client_version"` + LastStartTime string `json:"last_start_time"` + LastCloseTime string `json:"last_close_time"` + ClientAddr string `json:"client_addr"` +} + func GetNodeObservability(id uint, query NodeObservabilityQuery) (*NodeObservabilityView, error) { now := time.Now() node, err := model.GetNodeByID(id) @@ -90,7 +110,7 @@ func GetNodeObservability(id uint, query NodeObservabilityQuery) (*NodeObservabi return nil, err } - return &NodeObservabilityView{ + view := &NodeObservabilityView{ NodeID: node.NodeID, Profile: profile, MetricSnapshots: snapshots, @@ -107,7 +127,40 @@ func GetNodeObservability(id uint, query NodeObservabilityQuery) (*NodeObservabi Network24h: buildNetworkTrendPoints(now, trendSnapshots), DiskIO24h: buildDiskIOTrendPoints(now, trendSnapshots), }, - }, nil + } + if node.NodeType == "tunnel_relay" { + view.RelayDashboard = buildRelayDashboardSnapshot(node) + } + return view, nil +} + +func buildRelayDashboardSnapshot(node *model.Node) *RelayDashboardSnapshot { + if node == nil { + return nil + } + totalProxies := node.RelayFrpsProxyCount + if totalProxies < 0 { + totalProxies = 0 + } + onlineProxies := totalProxies + if node.RelayStatus != "healthy" { + onlineProxies = 0 + } + return &RelayDashboardSnapshot{ + TotalProxies: totalProxies, + OnlineProxies: onlineProxies, + OfflineProxies: totalProxies - onlineProxies, + Proxies: []RelayProxyStat{}, + TotalConnections: maxInt(node.RelayFrpsConnections, 0), + ClientCounts: 0, + } +} + +func maxInt(a int, b int) int { + if a > b { + return a + } + return b } func CleanupNodeHealthEvents(id uint) (*NodeHealthEventCleanupResult, error) { diff --git a/openflare_server/service/relay.go b/openflare_server/service/relay.go index 579e2296..c414fa12 100644 --- a/openflare_server/service/relay.go +++ b/openflare_server/service/relay.go @@ -11,11 +11,16 @@ import ( // RelayHeartbeatPayload is the payload sent by OpenFlareRelay in each heartbeat. type RelayHeartbeatPayload struct { - RelayVersion string `json:"relay_version"` - FrpVersion string `json:"frp_version"` - RelayStatus string `json:"relay_status"` - FrpsConnCount int `json:"frps_connections"` - FrpsProxyCount int `json:"frps_proxy_count"` + RelayVersion string `json:"relay_version"` + FrpVersion string `json:"frp_version"` + RelayStatus string `json:"relay_status"` + FrpsConnCount int `json:"frps_connections"` + FrpsProxyCount int `json:"frps_proxy_count"` + Name string `json:"name"` + IP string `json:"ip"` + Profile *AgentNodeSystemProfile `json:"profile,omitempty"` + Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"` + HealthEvents []AgentNodeHealthEvent `json:"health_events,omitempty"` } // RelayConfig is the frps configuration sent to the Relay. @@ -48,6 +53,8 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear payload.RelayVersion = strings.TrimSpace(payload.RelayVersion) payload.FrpVersion = strings.TrimSpace(payload.FrpVersion) payload.RelayStatus = normalizeRelayStatus(payload.RelayStatus) + payload.Name = strings.TrimSpace(payload.Name) + payload.IP = strings.TrimSpace(payload.IP) changes := make(map[string]any) appendRelayChange := func(key string, before any, after any) { @@ -59,6 +66,22 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear appendRelayChange("relay_version", node.RelayVersion, payload.RelayVersion) appendRelayChange("relay_frp_version", node.RelayFrpVersion, payload.FrpVersion) appendRelayChange("relay_status", node.RelayStatus, payload.RelayStatus) + appendRelayChange("relay_frps_connections", node.RelayFrpsConnections, payload.FrpsConnCount) + appendRelayChange("relay_frps_proxy_count", node.RelayFrpsProxyCount, payload.FrpsProxyCount) + if payload.Name != "" && strings.TrimSpace(node.Name) == "" { + appendRelayChange("name", node.Name, payload.Name) + node.Name = payload.Name + } + if payload.IP != "" && !node.IPManualOverride { + appendRelayChange("ip", node.IP, payload.IP) + node.IP = payload.IP + if !node.GeoManualOverride { + applyGeoInfoFromIP(node, node.IP) + changes["geo_name"] = node.GeoName + changes["geo_latitude"] = node.GeoLatitude + changes["geo_longitude"] = node.GeoLongitude + } + } if !node.LastSeenAt.Equal(now) { changes["last_seen_at"] = now } @@ -67,6 +90,8 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear node.RelayVersion = payload.RelayVersion node.RelayFrpVersion = payload.FrpVersion node.RelayStatus = payload.RelayStatus + node.RelayFrpsConnections = payload.FrpsConnCount + node.RelayFrpsProxyCount = payload.FrpsProxyCount node.LastSeenAt = now node.Status = NodeStatusOnline @@ -76,6 +101,7 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear } } refreshAgentTokenCache(node) + persistRelayHeartbeatObservability(node.NodeID, payload, node.LastSeenAt) return &RelayHeartbeatResponse{ RelayConfig: buildRelayConfig(node), @@ -83,6 +109,14 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear }, nil } +func persistRelayHeartbeatObservability(nodeID string, payload RelayHeartbeatPayload, reportedAt time.Time) { + persistHeartbeatObservability(nodeID, AgentNodePayload{ + Profile: payload.Profile, + Snapshot: payload.Snapshot, + HealthEvents: payload.HealthEvents, + }, reportedAt) +} + func buildRelayConfig(node *model.Node) *RelayConfig { if node == nil { return nil diff --git a/openflare_server/service/relay_test.go b/openflare_server/service/relay_test.go new file mode 100644 index 00000000..60cbb78d --- /dev/null +++ b/openflare_server/service/relay_test.go @@ -0,0 +1,98 @@ +package service + +import ( + "openflare/model" + "testing" + "time" +) + +func TestHeartbeatRelayPersistsRuntimeAndObservability(t *testing.T) { + setupServiceTestDB(t) + + node := &model.Node{ + NodeID: "node-relay-observe", + Name: "relay-1", + IP: "", + AgentToken: "relay-token", + Status: NodeStatusPending, + NodeType: "tunnel_relay", + RelayStatus: "unknown", + RelayVersion: "", + } + if err := node.Insert(); err != nil { + t.Fatalf("failed to seed relay node: %v", err) + } + + now := time.Now().UTC() + _, err := HeartbeatRelay(node, RelayHeartbeatPayload{ + RelayVersion: "v0.1.0", + FrpVersion: "0.61.0", + RelayStatus: "healthy", + FrpsConnCount: 7, + FrpsProxyCount: 3, + Name: "relay-runtime", + IP: "203.0.113.9", + Profile: &AgentNodeSystemProfile{ + Hostname: "relay-runtime", + OSName: "Ubuntu", + OSVersion: "24.04", + Architecture: "amd64", + CPUCores: 4, + ReportedAtUnix: now.Unix(), + }, + Snapshot: &AgentNodeMetricSnapshot{ + CapturedAtUnix: now.Unix(), + CPUUsagePercent: 12.5, + NetworkRxBytes: 1024, + NetworkTxBytes: 2048, + }, + HealthEvents: []AgentNodeHealthEvent{}, + }) + if err != nil { + t.Fatalf("HeartbeatRelay failed: %v", err) + } + + updated, err := model.GetNodeByNodeID(node.NodeID) + if err != nil { + t.Fatalf("failed to reload node: %v", err) + } + if updated.Status != NodeStatusOnline || updated.RelayStatus != "healthy" { + t.Fatalf("unexpected relay status: %+v", updated) + } + if updated.IP != "203.0.113.9" { + t.Fatalf("expected relay IP to be updated, got %q", updated.IP) + } + if updated.RelayVersion != "v0.1.0" || updated.RelayFrpVersion != "0.61.0" { + t.Fatalf("expected relay versions to be updated, got relay=%q frp=%q", updated.RelayVersion, updated.RelayFrpVersion) + } + if updated.RelayFrpsConnections != 7 || updated.RelayFrpsProxyCount != 3 { + t.Fatalf("expected relay counters to be stored, got connections=%d proxies=%d", updated.RelayFrpsConnections, updated.RelayFrpsProxyCount) + } + + profile, err := model.GetNodeSystemProfile(node.NodeID) + if err != nil { + t.Fatalf("expected relay system profile: %v", err) + } + if profile.Hostname != "relay-runtime" || profile.OSName != "Ubuntu" { + t.Fatalf("unexpected relay profile: %+v", profile) + } + + snapshots, err := model.ListNodeMetricSnapshots(node.NodeID, now.Add(-time.Minute), 10) + if err != nil { + t.Fatalf("failed to list relay snapshots: %v", err) + } + if len(snapshots) != 1 || snapshots[0].CPUUsagePercent != 12.5 { + t.Fatalf("unexpected relay snapshots: %+v", snapshots) + } + + observability, err := GetNodeObservability(updated.ID, NodeObservabilityQuery{Hours: 1, Limit: 10}) + if err != nil { + t.Fatalf("GetNodeObservability failed: %v", err) + } + if observability.RelayDashboard == nil { + t.Fatal("expected relay dashboard snapshot") + } + if observability.RelayDashboard.TotalConnections != 7 || observability.RelayDashboard.TotalProxies != 3 { + t.Fatalf("unexpected relay dashboard: %+v", observability.RelayDashboard) + } +} diff --git a/openflare_server/web/features/nodes/components/node-detail-page.tsx b/openflare_server/web/features/nodes/components/node-detail-page.tsx index cf2c9cdf..56633ceb 100644 --- a/openflare_server/web/features/nodes/components/node-detail-page.tsx +++ b/openflare_server/web/features/nodes/components/node-detail-page.tsx @@ -1252,7 +1252,7 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { {node.node_type === 'tunnel_relay' ? (
| - {node.agent_version || 'unknown'} + {node.node_type === 'tunnel_relay' + ? node.relay_version || node.agent_version || 'unknown' + : node.agent_version || 'unknown'} |
diff --git a/openflare_server/web/features/nodes/types.ts b/openflare_server/web/features/nodes/types.ts
index 6caf3022..974c484f 100644
--- a/openflare_server/web/features/nodes/types.ts
+++ b/openflare_server/web/features/nodes/types.ts
@@ -13,6 +13,10 @@ export interface NodeItem {
relay_client_proxy_url: string;
relay_auth_token: string;
relay_status: string;
+ relay_frp_version: string;
+ relay_version: string;
+ relay_frps_connections: number;
+ relay_frps_proxy_count: number;
geo_name: string;
geo_latitude?: number | null;
geo_longitude?: number | null;
diff --git a/openflared/Dockerfile b/openflared/Dockerfile
index 78249ce1..916a9bb0 100644
--- a/openflared/Dockerfile
+++ b/openflared/Dockerfile
@@ -1,5 +1,9 @@
+ARG VERSION=dev
+
FROM golang:1.25-alpine AS builder
+ARG VERSION
+
WORKDIR /build
COPY openflared/go.mod openflared/go.sum ./
@@ -8,7 +12,7 @@ RUN go mod download
COPY openflare_server/ ../openflare_server/
COPY openflared/ .
-RUN CGO_ENABLED=0 GOOS=linux go build -o flared ./cmd/flared
+RUN CGO_ENABLED=0 GOOS=linux go build -trimpath -ldflags "-s -w -X 'openflare-flared/internal/config.Version=$VERSION'" -o flared ./cmd/flared
# Final runtime image
FROM fatedier/frpc:v0.69.0
diff --git a/openflared/internal/config/version.go b/openflared/internal/config/version.go
new file mode 100644
index 00000000..f356490f
--- /dev/null
+++ b/openflared/internal/config/version.go
@@ -0,0 +1,3 @@
+package config
+
+var Version = "dev"
diff --git a/openflared/internal/heartbeat/service.go b/openflared/internal/heartbeat/service.go
index 1ae0ec5c..a8907f2b 100644
--- a/openflared/internal/heartbeat/service.go
+++ b/openflared/internal/heartbeat/service.go
@@ -46,7 +46,7 @@ func (s *Service) doHeartbeat(ctx context.Context) {
slog.Debug("sending flared heartbeat")
payload := service.FlaredHeartbeatPayload{
- ClientVersion: "0.1.0", // TODO dynamically inject build version
+ ClientVersion: config.Version,
FrpVersion: s.frpcManager.GetVersion(),
TunnelStatus: "running", // TODO implement proper status tracking
ConnectedRelays: s.frpcManager.GetConnectedRelays(),
|