diff --git a/openflare_agent/cmd/agent/main.go b/openflare_agent/cmd/agent/main.go index 290cd9aa..90d32f82 100644 --- a/openflare_agent/cmd/agent/main.go +++ b/openflare_agent/cmd/agent/main.go @@ -32,7 +32,7 @@ func main() { slog.Error("load agent config failed", "error", err) os.Exit(1) } - cfg.NginxVersion = nginx.DetectVersion( + cfg.ExtVersion = nginx.DetectVersion( context.Background(), nginx.ExecutorOptions{ NginxPath: cfg.OpenrestyPath, diff --git a/openflare_agent/internal/agent/runner.go b/openflare_agent/internal/agent/runner.go index e5e80f25..598c93d3 100644 --- a/openflare_agent/internal/agent/runner.go +++ b/openflare_agent/internal/agent/runner.go @@ -74,7 +74,7 @@ func (r *Runner) Run(ctx context.Context) error { return err } slog.Info("agent runner started", "node_id", nodeID, "node", r.Config.NodeName, "ip", r.Config.NodeIP) - if r.hasAgentToken() { + if r.hasAccessToken() { if _, hbErr := r.performHeartbeatCycle(ctx, nodeID, true); hbErr != nil { slog.Error("agent startup heartbeat failed", "error", hbErr) } @@ -119,7 +119,7 @@ func (r *Runner) Run(ctx context.Context) error { delay := wsBackoff.Next() nextWSAttempt = time.Now().Add(delay) slog.Debug("agent ws disconnected; resuming http heartbeat", "retry_after", delay, "error", wsErr) - if r.hasAgentToken() { + if r.hasAccessToken() { if _, hbErr := r.performHeartbeatCycle(ctx, nodeID, false); hbErr != nil { slog.Error("agent heartbeat after ws disconnect failed", "error", hbErr) } @@ -128,7 +128,7 @@ func (r *Runner) Run(ctx context.Context) error { if wsDone != nil { continue } - if !r.hasAgentToken() { + if !r.hasAccessToken() { if err = r.tryRegister(ctx, &nodeID); err != nil { slog.Error("agent discovery register failed", "error", err) } @@ -181,7 +181,7 @@ func (r *Runner) performHeartbeatCycle(ctx context.Context, nodeID string, start } func (r *Runner) shouldUseWebSocket() bool { - enabled := r.WebSocketService != nil && r.websocketUpgradeEnabled && r.hasAgentToken() + enabled := r.WebSocketService != nil && r.websocketUpgradeEnabled && r.hasAccessToken() slog.Debug("agent ws upgrade eligibility checked", "enabled", enabled, "server_enabled", r.websocketUpgradeEnabled, "url", r.websocketURL()) return enabled } @@ -365,8 +365,8 @@ func (backoff *webSocketBackoff) Reset() { } } -func (r *Runner) hasAgentToken() bool { - return strings.TrimSpace(r.Config.AgentToken) != "" +func (r *Runner) hasAccessToken() bool { + return strings.TrimSpace(r.Config.AccessToken) != "" } func (r *Runner) applySettings(settings *protocol.AgentSettings) bool { @@ -447,7 +447,7 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error { if err != nil { return err } - if response == nil || strings.TrimSpace(response.AgentToken) == "" || strings.TrimSpace(response.NodeID) == "" { + if response == nil || strings.TrimSpace(response.AccessToken) == "" || strings.TrimSpace(response.NodeID) == "" { return errors.New("discovery register response 缺少 node_id 或 agent_token") } snapshot, err := r.StateStore.Load() @@ -458,14 +458,14 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error { if err = r.StateStore.Save(snapshot); err != nil { return err } - r.Config.AgentToken = response.AgentToken + r.Config.AccessToken = response.AccessToken r.Config.DiscoveryToken = "" if err = r.Config.Save(); err != nil { return err } - r.HeartbeatService.SetToken(response.AgentToken) + r.HeartbeatService.SetToken(response.AccessToken) if r.WebSocketService != nil { - r.WebSocketService.SetToken(response.AgentToken) + r.WebSocketService.SetToken(response.AccessToken) } *nodeID = response.NodeID slog.Info("agent discovery registration succeeded", "node_id", response.NodeID) @@ -573,23 +573,25 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { if managedOpenRestyMetrics == nil { managedOpenRestyMetrics = fallbackMetrics } - metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore, managedOpenRestyMetrics) + metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore) + openrestyObservation := observability.BuildOpenrestyObservation(managedOpenRestyMetrics) healthEvents := observability.BuildHealthEvents(snapshot) payload := protocol.NodePayload{ - NodeID: nodeID, - Name: r.Config.NodeName, - IP: r.Config.NodeIP, - AgentVersion: r.Config.AgentVersion, - NginxVersion: r.Config.NginxVersion, - CurrentVersion: snapshot.CurrentVersion, - LastError: snapshot.LastError, - OpenrestyStatus: openrestyStatus, - OpenrestyMessage: snapshot.OpenrestyMessage, - Profile: profile, - Snapshot: metricSnapshot, - TrafficReport: trafficReport, - AccessLogs: accessLogs, - HealthEvents: healthEvents, + NodeID: nodeID, + Name: r.Config.NodeName, + IP: r.Config.NodeIP, + Version: r.Config.Version, + ExtVersion: r.Config.ExtVersion, + CurrentVersion: snapshot.CurrentVersion, + LastError: snapshot.LastError, + OpenrestyStatus: openrestyStatus, + OpenrestyMessage: snapshot.OpenrestyMessage, + Profile: profile, + Snapshot: metricSnapshot, + OpenrestyObservation: openrestyObservation, + TrafficReport: trafficReport, + AccessLogs: accessLogs, + HealthEvents: healthEvents, } if r.SyncService != nil { checksums, err := r.SyncService.WAFIPGroupChecksums() @@ -619,17 +621,18 @@ func (r *Runner) prepareHeartbeatPayload(nodeID string) (protocol.NodePayload, [ } now := time.Now().UTC() retainAfterUnix := now.Add(-time.Duration(r.Config.ObservabilityReplayMinutes) * time.Minute).Unix() - windowStartedAtUnix := state.ObservabilityWindowStartedAt(payload.Snapshot, payload.TrafficReport) + windowStartedAtUnix := state.ObservabilityWindowStartedAt(payload.Snapshot, payload.OpenrestyObservation, payload.TrafficReport) if windowStartedAtUnix <= 0 { return payload, nil } record := state.ObservabilityBufferRecord{ - WindowStartedAtUnix: windowStartedAtUnix, - Snapshot: payload.Snapshot, - TrafficReport: payload.TrafficReport, - AccessLogs: payload.AccessLogs, - QueuedAtUnix: now.Unix(), + WindowStartedAtUnix: windowStartedAtUnix, + Snapshot: payload.Snapshot, + OpenrestyObservation: payload.OpenrestyObservation, + TrafficReport: payload.TrafficReport, + AccessLogs: payload.AccessLogs, + QueuedAtUnix: now.Unix(), } if err := r.ObservabilityBuffer.Upsert(record, retainAfterUnix); err != nil { slog.Error("upsert observability buffer failed", "error", err) @@ -649,10 +652,11 @@ func (r *Runner) prepareHeartbeatPayload(nodeID string) (protocol.NodePayload, [ continue } buffered = append(buffered, protocol.BufferedObservabilityRecord{ - WindowStartedAtUnix: item.WindowStartedAtUnix, - Snapshot: item.Snapshot, - TrafficReport: item.TrafficReport, - AccessLogs: item.AccessLogs, + WindowStartedAtUnix: item.WindowStartedAtUnix, + Snapshot: item.Snapshot, + OpenrestyObservation: item.OpenrestyObservation, + TrafficReport: item.TrafficReport, + AccessLogs: item.AccessLogs, }) ackWindows = append(ackWindows, item.WindowStartedAtUnix) } diff --git a/openflare_agent/internal/agent/runner_test.go b/openflare_agent/internal/agent/runner_test.go index 168d3067..45f4ca14 100644 --- a/openflare_agent/internal/agent/runner_test.go +++ b/openflare_agent/internal/agent/runner_test.go @@ -194,11 +194,11 @@ func TestRunnerKeepsHeartbeatWhenStartupSyncFails(t *testing.T) { } runner := &Runner{ Config: &config.Config{ - AgentToken: "agent-token", + AccessToken: "agent-token", NodeName: "edge-01", NodeIP: "10.0.0.8", - AgentVersion: config.AgentVersion, - NginxVersion: "1.27.1.2", + Version: config.Version, + ExtVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), }, StateStore: stateStore, @@ -247,11 +247,11 @@ func TestRunnerDoesNotExitOnHeartbeatOrSyncError(t *testing.T) { } runner := &Runner{ Config: &config.Config{ - AgentToken: "agent-token", + AccessToken: "agent-token", NodeName: "edge-01", NodeIP: "10.0.0.8", - AgentVersion: config.AgentVersion, - NginxVersion: "1.27.1.2", + Version: config.Version, + ExtVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), }, StateStore: stateStore, @@ -305,11 +305,11 @@ func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) { } runner := &Runner{ Config: &config.Config{ - AgentToken: "agent-token", + AccessToken: "agent-token", NodeName: "edge-01", NodeIP: "10.0.0.8", - AgentVersion: config.AgentVersion, - NginxVersion: "1.27.1.2", + Version: config.Version, + ExtVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), }, StateStore: stateStore, @@ -361,8 +361,8 @@ func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) { Config: &config.Config{ NodeName: "edge-observe-1", NodeIP: "10.0.0.51", - AgentVersion: config.AgentVersion, - NginxVersion: "1.27.1.2", + Version: config.Version, + ExtVersion: "1.27.1.2", DataDir: tempDir, RouteConfigPath: filepath.Join(tempDir, "conf.d", "openflare_routes.conf"), AccessLogPath: filepath.Join(tempDir, "var", "log", "openflare", "access.log"), @@ -441,11 +441,11 @@ func TestRunnerReplaysBufferedObservabilityAfterHeartbeatRecovery(t *testing.T) } runner := &Runner{ Config: &config.Config{ - AgentToken: "agent-token", + AccessToken: "agent-token", NodeName: "edge-buffer-01", NodeIP: "10.0.0.52", - AgentVersion: config.AgentVersion, - NginxVersion: "1.27.1.2", + Version: config.Version, + ExtVersion: "1.27.1.2", DataDir: tempDir, RouteConfigPath: filepath.Join(tempDir, "conf.d", "openflare_routes.conf"), HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), @@ -498,9 +498,9 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) heartbeatService := &fakeHeartbeatService{ registerResp: &protocol.RegisterNodeResponse{ - NodeID: "node-server-assigned", - AgentToken: "agent-token-issued", - Name: "edge-01", + NodeID: "node-server-assigned", + AccessToken: "agent-token-issued", + Name: "edge-01", }, heartbeatResults: []*protocol.HeartbeatResult{{}}, onHeartbeat: func(callCount int) { @@ -524,8 +524,8 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { DiscoveryToken: cfg.DiscoveryToken, NodeName: cfg.NodeName, NodeIP: cfg.NodeIP, - AgentVersion: config.AgentVersion, - NginxVersion: "1.27.1.2", + Version: config.Version, + ExtVersion: "1.27.1.2", HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), }, StateStore: stateStore, @@ -533,8 +533,8 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { SyncService: syncService, } runner.Config = cfg - runner.Config.AgentVersion = config.AgentVersion - runner.Config.NginxVersion = "1.27.1.2" + runner.Config.Version = config.Version + runner.Config.ExtVersion = "1.27.1.2" runner.Config.HeartbeatInterval = config.MillisecondDuration(10 * time.Millisecond) err = runner.Run(ctx) @@ -554,7 +554,7 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { if snapshot.NodeID != "node-server-assigned" { t.Fatalf("expected node id to be replaced, got %q", snapshot.NodeID) } - if runner.Config.AgentToken != "agent-token-issued" || runner.Config.DiscoveryToken != "" { + if runner.Config.AccessToken != "agent-token-issued" || runner.Config.DiscoveryToken != "" { t.Fatal("expected config token rotation to complete") } } diff --git a/openflare_agent/internal/config/config.go b/openflare_agent/internal/config/config.go index a5a0d399..280cbf45 100644 --- a/openflare_agent/internal/config/config.go +++ b/openflare_agent/internal/config/config.go @@ -40,12 +40,12 @@ var ( type Config struct { ServerURL string `json:"server_url"` - AgentToken string `json:"agent_token"` + AccessToken string `json:"agent_token"` DiscoveryToken string `json:"discovery_token"` NodeName string `json:"node_name"` NodeIP string `json:"node_ip"` - AgentVersion string `json:"-"` - NginxVersion string `json:"-"` + Version string `json:"-"` + ExtVersion string `json:"-"` OpenrestyPath string `json:"openresty_path"` OpenrestyResolvers []string `json:"openresty_resolvers,omitempty"` DataDir string `json:"data_dir"` @@ -71,7 +71,7 @@ type Config struct { type configFile struct { ServerURL string `json:"server_url"` - AgentToken string `json:"agent_token"` + AccessToken string `json:"agent_token"` DiscoveryToken string `json:"discovery_token"` NodeName string `json:"node_name"` NodeIP string `json:"node_ip"` @@ -113,7 +113,7 @@ func Load(path string) (*Config, error) { } cfg := &Config{ ServerURL: file.ServerURL, - AgentToken: file.AgentToken, + AccessToken: file.AccessToken, DiscoveryToken: file.DiscoveryToken, NodeName: file.NodeName, NodeIP: file.NodeIP, @@ -149,7 +149,7 @@ func Load(path string) (*Config, error) { func applyDefaults(cfg *Config, baseDir string) { baseDir = filepath.Clean(baseDir) - cfg.AgentVersion = AgentVersion + cfg.Version = Version cfg.OpenrestyResolvers = utils.UniqueAndCleanStringSlice(cfg.OpenrestyResolvers) if cfg.OpenrestyPath == "" { cfg.OpenrestyPath = "openresty" @@ -275,7 +275,7 @@ func applyEnvOverrides(cfg *Config) { } } overrideString("OPENFLARE_SERVER_URL", &cfg.ServerURL) - overrideString("OPENFLARE_AGENT_TOKEN", &cfg.AgentToken) + overrideString("OPENFLARE_AGENT_TOKEN", &cfg.AccessToken) overrideString("OPENFLARE_DISCOVERY_TOKEN", &cfg.DiscoveryToken) overrideString("OPENFLARE_NODE_NAME", &cfg.NodeName) overrideString("OPENFLARE_NODE_IP", &cfg.NodeIP) @@ -336,7 +336,7 @@ func validate(cfg *Config) error { if cfg.ServerURL == "" { return errors.New("server_url 不能为空") } - if strings.TrimSpace(cfg.AgentToken) == "" && strings.TrimSpace(cfg.DiscoveryToken) == "" { + if strings.TrimSpace(cfg.AccessToken) == "" && strings.TrimSpace(cfg.DiscoveryToken) == "" { return errors.New("agent_token 和 discovery_token 不能同时为空") } if cfg.NodeName == "" { @@ -361,7 +361,7 @@ func (cfg *Config) InitialAuthToken() string { if cfg == nil { return "" } - if token := strings.TrimSpace(cfg.AgentToken); token != "" { + if token := strings.TrimSpace(cfg.AccessToken); token != "" { return token } return strings.TrimSpace(cfg.DiscoveryToken) diff --git a/openflare_agent/internal/config/config_test.go b/openflare_agent/internal/config/config_test.go index 2573591f..6328035e 100644 --- a/openflare_agent/internal/config/config_test.go +++ b/openflare_agent/internal/config/config_test.go @@ -226,7 +226,7 @@ func TestLoadUsesEnvConfigWhenFileIsMissing(t *testing.T) { if err != nil { t.Fatalf("Load failed: %v", err) } - if cfg.ServerURL != "http://127.0.0.1:3000" || cfg.AgentToken != "token" { + if cfg.ServerURL != "http://127.0.0.1:3000" || cfg.AccessToken != "token" { t.Fatalf("unexpected env auth config: %#v", cfg) } if cfg.OpenrestyPath != "/usr/bin/openresty" { @@ -334,8 +334,8 @@ func TestLoadEnvOverridesConfigFile(t *testing.T) { if cfg.ServerURL != "http://new:3000" { t.Fatalf("expected server url from env, got %s", cfg.ServerURL) } - if cfg.AgentToken != "new-token" { - t.Fatalf("expected token from env, got %s", cfg.AgentToken) + if cfg.AccessToken != "new-token" { + t.Fatalf("expected token from env, got %s", cfg.AccessToken) } if cfg.OpenrestyPath != "/new/openresty" { t.Fatalf("expected openresty path from env, got %s", cfg.OpenrestyPath) @@ -385,7 +385,7 @@ func TestSavePersistsMillisecondsAndOmitsRuntimeVersions(t *testing.T) { if err != nil { t.Fatalf("Load failed: %v", err) } - cfg.NginxVersion = "1.27.1.2" + cfg.ExtVersion = "1.27.1.2" cfg.HeartbeatInterval = MillisecondDuration(5 * time.Second) cfg.RequestTimeout = MillisecondDuration(7 * time.Second) cfg.OpenrestyResolvers = []string{"10.0.0.2", "1.1.1.1"} @@ -461,7 +461,7 @@ func TestInitialAuthToken(t *testing.T) { var cfg *Config if tt.name != "nil config returns empty string" { cfg = &Config{ - AgentToken: tt.agentToken, + AccessToken: tt.agentToken, DiscoveryToken: tt.discoveryToken, } } diff --git a/openflare_agent/internal/config/version.go b/openflare_agent/internal/config/version.go index 89f7ad9b..f356490f 100644 --- a/openflare_agent/internal/config/version.go +++ b/openflare_agent/internal/config/version.go @@ -1,3 +1,3 @@ package config -var AgentVersion = "dev" +var Version = "dev" diff --git a/openflare_agent/internal/nginx/manager.go b/openflare_agent/internal/nginx/manager.go index e67745dc..29fb30c5 100644 --- a/openflare_agent/internal/nginx/manager.go +++ b/openflare_agent/internal/nginx/manager.go @@ -541,7 +541,7 @@ func detectVersion(ctx context.Context, options ExecutorOptions, runner CommandR if err != nil { return "", fmt.Errorf("run runtime -v failed: %w: %s", err, string(output)) } - version := parseNginxVersion(string(output)) + version := parseExtVersion(string(output)) if version == "" { return "", errors.New("cannot parse runtime version from binary output") } @@ -550,7 +550,7 @@ func detectVersion(ctx context.Context, options ExecutorOptions, runner CommandR return "", errors.New("openresty path is empty") } -func parseNginxVersion(output string) string { +func parseExtVersion(output string) string { matches := nginxVersionPattern.FindStringSubmatch(output) if len(matches) != 2 { return "" diff --git a/openflare_agent/internal/nginx/manager_test.go b/openflare_agent/internal/nginx/manager_test.go index 4d18fe2f..7d20da6e 100644 --- a/openflare_agent/internal/nginx/manager_test.go +++ b/openflare_agent/internal/nginx/manager_test.go @@ -256,13 +256,13 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) { } } -func TestParseNginxVersionIgnoresDockerEntrypointPaths(t *testing.T) { +func TestParseExtVersionIgnoresDockerEntrypointPaths(t *testing.T) { output := strings.Join([]string{ "/docker-entrypoint.sh: /docker-entrypoint.d/10-listen-on-ipv6-by-default.sh: info: can not modify /etc/nginx/conf.d/default.conf (read-only file system?)", "nginx version: openresty/1.27.1.2", }, "\n") - version := parseNginxVersion(output) + version := parseExtVersion(output) if version != "1.27.1.2" { t.Fatalf("unexpected version: %s", version) } diff --git a/openflare_agent/internal/observability/collector.go b/openflare_agent/internal/observability/collector.go index a5f3fdad..54a5b61a 100644 --- a/openflare_agent/internal/observability/collector.go +++ b/openflare_agent/internal/observability/collector.go @@ -40,7 +40,7 @@ func BuildProfile(cfg *config.Config, stateStore *state.Store) *protocol.NodeSys return profile } -func BuildSnapshot(cfg *config.Config, stateStore *state.Store, managed *ManagedOpenRestyMetrics) *protocol.NodeMetricSnapshot { +func BuildSnapshot(cfg *config.Config, stateStore *state.Store) *protocol.NodeMetricSnapshot { now := time.Now().UTC() metric := &protocol.NodeMetricSnapshot{ CapturedAtUnix: now.Unix(), @@ -56,11 +56,6 @@ func BuildSnapshot(cfg *config.Config, stateStore *state.Store, managed *Managed metric.NetworkRxBytes, metric.NetworkTxBytes = readLinuxNetworkTotals() metric.DiskReadBytes, metric.DiskWriteBytes = readLinuxDiskTotals() - if managed != nil { - metric.OpenrestyRxBytes = managed.OpenrestyRxBytes - metric.OpenrestyTxBytes = managed.OpenrestyTxBytes - metric.OpenrestyConnections = managed.OpenrestyConnections - } if stateStore == nil { return metric @@ -86,6 +81,18 @@ func BuildSnapshot(cfg *config.Config, stateStore *state.Store, managed *Managed return metric } +func BuildOpenrestyObservation(managed *ManagedOpenRestyMetrics) *protocol.NodeOpenrestyObservation { + if managed == nil { + return nil + } + return &protocol.NodeOpenrestyObservation{ + CapturedAtUnix: time.Now().UTC().Unix(), + OpenrestyRxBytes: managed.OpenrestyRxBytes, + OpenrestyTxBytes: managed.OpenrestyTxBytes, + OpenrestyConnections: managed.OpenrestyConnections, + } +} + func BuildHealthEvents(snapshot *state.Snapshot) []protocol.NodeHealthEvent { if snapshot == nil { return []protocol.NodeHealthEvent{} diff --git a/openflare_agent/internal/protocol/agent_api.go b/openflare_agent/internal/protocol/agent_api.go index 70472632..5a860701 100644 --- a/openflare_agent/internal/protocol/agent_api.go +++ b/openflare_agent/internal/protocol/agent_api.go @@ -72,14 +72,15 @@ type NodePayload struct { NodeID string `json:"node_id"` Name string `json:"name"` IP string `json:"ip"` - AgentVersion string `json:"agent_version"` - NginxVersion string `json:"nginx_version"` + Version string `json:"agent_version"` + ExtVersion string `json:"nginx_version"` CurrentVersion string `json:"current_version"` LastError string `json:"last_error"` OpenrestyStatus string `json:"openresty_status"` OpenrestyMessage string `json:"openresty_message"` Profile *NodeSystemProfile `json:"profile,omitempty"` Snapshot *NodeMetricSnapshot `json:"snapshot,omitempty"` + OpenrestyObservation *NodeOpenrestyObservation `json:"openresty_observation,omitempty"` TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"` AccessLogs []NodeAccessLog `json:"access_logs,omitempty"` BufferedObservability []BufferedObservabilityRecord `json:"buffered_observability,omitempty"` @@ -102,19 +103,23 @@ type NodeSystemProfile struct { } type NodeMetricSnapshot struct { - CapturedAtUnix int64 `json:"captured_at_unix"` - CPUUsagePercent float64 `json:"cpu_usage_percent"` - MemoryUsedBytes int64 `json:"memory_used_bytes"` - MemoryTotalBytes int64 `json:"memory_total_bytes"` - StorageUsedBytes int64 `json:"storage_used_bytes"` - StorageTotalBytes int64 `json:"storage_total_bytes"` - DiskReadBytes int64 `json:"disk_read_bytes"` - DiskWriteBytes int64 `json:"disk_write_bytes"` - NetworkRxBytes int64 `json:"network_rx_bytes"` - NetworkTxBytes int64 `json:"network_tx_bytes"` - OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` - OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` - OpenrestyConnections int64 `json:"openresty_connections"` + CapturedAtUnix int64 `json:"captured_at_unix"` + CPUUsagePercent float64 `json:"cpu_usage_percent"` + MemoryUsedBytes int64 `json:"memory_used_bytes"` + MemoryTotalBytes int64 `json:"memory_total_bytes"` + StorageUsedBytes int64 `json:"storage_used_bytes"` + StorageTotalBytes int64 `json:"storage_total_bytes"` + DiskReadBytes int64 `json:"disk_read_bytes"` + DiskWriteBytes int64 `json:"disk_write_bytes"` + NetworkRxBytes int64 `json:"network_rx_bytes"` + NetworkTxBytes int64 `json:"network_tx_bytes"` +} + +type NodeOpenrestyObservation struct { + CapturedAtUnix int64 `json:"captured_at_unix"` + OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` + OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` + OpenrestyConnections int64 `json:"openresty_connections"` } type NodeTrafficReport struct { @@ -137,10 +142,11 @@ type NodeAccessLog struct { } type BufferedObservabilityRecord struct { - WindowStartedAtUnix int64 `json:"window_started_at_unix"` - Snapshot *NodeMetricSnapshot `json:"snapshot,omitempty"` - TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"` - AccessLogs []NodeAccessLog `json:"access_logs,omitempty"` + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + Snapshot *NodeMetricSnapshot `json:"snapshot,omitempty"` + OpenrestyObservation *NodeOpenrestyObservation `json:"openresty_observation,omitempty"` + TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"` + AccessLogs []NodeAccessLog `json:"access_logs,omitempty"` } type NodeHealthEvent struct { @@ -152,9 +158,9 @@ type NodeHealthEvent struct { } type RegisterNodeResponse struct { - NodeID string `json:"node_id"` - AgentToken string `json:"agent_token"` - Name string `json:"name"` + NodeID string `json:"node_id"` + AccessToken string `json:"agent_token"` + Name string `json:"name"` } type ApplyLogPayload struct { diff --git a/openflare_agent/internal/state/observability_buffer.go b/openflare_agent/internal/state/observability_buffer.go index ebf63f35..236564db 100644 --- a/openflare_agent/internal/state/observability_buffer.go +++ b/openflare_agent/internal/state/observability_buffer.go @@ -14,11 +14,12 @@ import ( const observabilityBufferWindowSeconds = 60 type ObservabilityBufferRecord struct { - WindowStartedAtUnix int64 `json:"window_started_at_unix"` - Snapshot *protocol.NodeMetricSnapshot `json:"snapshot,omitempty"` - TrafficReport *protocol.NodeTrafficReport `json:"traffic_report,omitempty"` - AccessLogs []protocol.NodeAccessLog `json:"access_logs,omitempty"` - QueuedAtUnix int64 `json:"queued_at_unix"` + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + Snapshot *protocol.NodeMetricSnapshot `json:"snapshot,omitempty"` + OpenrestyObservation *protocol.NodeOpenrestyObservation `json:"openresty_observation,omitempty"` + TrafficReport *protocol.NodeTrafficReport `json:"traffic_report,omitempty"` + AccessLogs []protocol.NodeAccessLog `json:"access_logs,omitempty"` + QueuedAtUnix int64 `json:"queued_at_unix"` } type ObservabilityBufferStore struct { @@ -31,7 +32,7 @@ func NewObservabilityBufferStore(path string) *ObservabilityBufferStore { } func (s *ObservabilityBufferStore) Upsert(record ObservabilityBufferRecord, retainAfterUnix int64) error { - if s == nil || record.WindowStartedAtUnix <= 0 || (record.Snapshot == nil && record.TrafficReport == nil && len(record.AccessLogs) == 0) { + if s == nil || record.WindowStartedAtUnix <= 0 || (record.Snapshot == nil && record.OpenrestyObservation == nil && record.TrafficReport == nil && len(record.AccessLogs) == 0) { return nil } s.mu.Lock() @@ -65,6 +66,9 @@ func mergeObservabilityBufferRecord(existing ObservabilityBufferRecord, incoming if incoming.Snapshot != nil { merged.Snapshot = incoming.Snapshot } + if incoming.OpenrestyObservation != nil { + merged.OpenrestyObservation = incoming.OpenrestyObservation + } if incoming.TrafficReport != nil { merged.TrafficReport = incoming.TrafficReport } @@ -191,10 +195,13 @@ func (s *ObservabilityBufferStore) saveUnlocked(records []ObservabilityBufferRec return os.WriteFile(s.path, data, 0o644) } -func ObservabilityWindowStartedAt(snapshot *protocol.NodeMetricSnapshot, traffic *protocol.NodeTrafficReport) int64 { +func ObservabilityWindowStartedAt(snapshot *protocol.NodeMetricSnapshot, openresty *protocol.NodeOpenrestyObservation, traffic *protocol.NodeTrafficReport) int64 { if traffic != nil && traffic.WindowStartedAtUnix > 0 { return traffic.WindowStartedAtUnix - (traffic.WindowStartedAtUnix % observabilityBufferWindowSeconds) } + if openresty != nil && openresty.CapturedAtUnix > 0 { + return openresty.CapturedAtUnix - (openresty.CapturedAtUnix % observabilityBufferWindowSeconds) + } if snapshot == nil || snapshot.CapturedAtUnix <= 0 { return 0 } diff --git a/openflare_agent/internal/state/observability_buffer_test.go b/openflare_agent/internal/state/observability_buffer_test.go index 19859606..93aab010 100644 --- a/openflare_agent/internal/state/observability_buffer_test.go +++ b/openflare_agent/internal/state/observability_buffer_test.go @@ -89,10 +89,10 @@ func TestObservabilityBufferStoreMergesAccessLogsWithinWindow(t *testing.T) { } func TestObservabilityWindowStartedAt(t *testing.T) { - if value := ObservabilityWindowStartedAt(nil, &protocol.NodeTrafficReport{WindowStartedAtUnix: 1710403200}); value != 1710403200 { + if value := ObservabilityWindowStartedAt(nil, nil, &protocol.NodeTrafficReport{WindowStartedAtUnix: 1710403200}); value != 1710403200 { t.Fatalf("unexpected traffic window start: %d", value) } - if value := ObservabilityWindowStartedAt(&protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403259}, nil); value != 1710403200 { + if value := ObservabilityWindowStartedAt(&protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403259}, nil, nil); value != 1710403200 { t.Fatalf("unexpected snapshot-derived window start: %d", value) } } diff --git a/openflare_agent/internal/updater/updater.go b/openflare_agent/internal/updater/updater.go index 46ab41d1..44986180 100644 --- a/openflare_agent/internal/updater/updater.go +++ b/openflare_agent/internal/updater/updater.go @@ -56,7 +56,7 @@ func (s *Service) CheckAndUpdate(ctx context.Context, repo string, options agent } remoteVersion := normalizeVersion(release.TagName) - localVersion := normalizeVersion(config.AgentVersion) + localVersion := normalizeVersion(config.Version) checkKey := buildReleaseCheckKey(options, remoteVersion) if remoteVersion == localVersion { diff --git a/openflare_agent/internal/updater/updater_test.go b/openflare_agent/internal/updater/updater_test.go index a66c598b..f4a34293 100644 --- a/openflare_agent/internal/updater/updater_test.go +++ b/openflare_agent/internal/updater/updater_test.go @@ -75,10 +75,10 @@ func TestGetReleaseByTag(t *testing.T) { } func TestCheckAndUpdateRequiresChecksumAsset(t *testing.T) { - originalVersion := config.AgentVersion - config.AgentVersion = "v1.0.0" + originalVersion := config.Version + config.Version = "v1.0.0" t.Cleanup(func() { - config.AgentVersion = originalVersion + config.Version = originalVersion }) assetName := assetNameForGOOSGOARCH(runtime.GOOS, runtime.GOARCH) diff --git a/openflare_relay/internal/heartbeat/service.go b/openflare_relay/internal/heartbeat/service.go index c8cd24a1..3d93645e 100644 --- a/openflare_relay/internal/heartbeat/service.go +++ b/openflare_relay/internal/heartbeat/service.go @@ -51,8 +51,8 @@ func (s *Service) doHeartbeat(ctx context.Context) { runtimeStatus := s.frpsManager.GetRuntimeStatus() payload := service.RelayHeartbeatPayload{ - RelayVersion: config.Version, - FrpVersion: s.frpsManager.GetVersion(), + Version: config.Version, + ExtVersion: s.frpsManager.GetVersion(), RelayStatus: runtimeStatus.Status, FrpsConnCount: runtimeStatus.Connections, FrpsProxyCount: runtimeStatus.ProxyCount, diff --git a/openflare_server/common/constants.go b/openflare_server/common/constants.go index ece6f3c1..be663d74 100644 --- a/openflare_server/common/constants.go +++ b/openflare_server/common/constants.go @@ -44,7 +44,7 @@ var WeChatServerAddress = "" var WeChatServerToken = "" var WeChatAccountQRCodeImageURL = "" -var AgentToken = "" +var AccessToken = "" var AgentDiscoveryToken = "" var NodeOfflineThreshold = 2 * time.Minute diff --git a/openflare_server/common/init.go b/openflare_server/common/init.go index a28038d9..ff932412 100644 --- a/openflare_server/common/init.go +++ b/openflare_server/common/init.go @@ -59,7 +59,7 @@ func ParseFlags() { } if os.Getenv("AGENT_TOKEN") != "" { - AgentToken = os.Getenv("AGENT_TOKEN") + AccessToken = os.Getenv("AGENT_TOKEN") } SetLogLevel(os.Getenv("LOG_LEVEL")) if *LogDir != "" { diff --git a/openflare_server/controller/agent.go b/openflare_server/controller/agent.go index c1ef29eb..5c9b45cf 100644 --- a/openflare_server/controller/agent.go +++ b/openflare_server/controller/agent.go @@ -19,7 +19,7 @@ import ( // @Tags Agent // @Accept json // @Produce json -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Param payload body service.AgentNodePayload true "Agent node payload" // @Success 200 {object} map[string]interface{} // @Failure 400 {object} map[string]interface{} @@ -36,7 +36,7 @@ func AgentRegister(c *gin.Context) { err error ) if authNode, ok := c.Get("agent_node"); ok { - result, err = service.RegisterNodeWithAgentToken(authNode.(*model.Node), payload) + result, err = service.RegisterNodeWithAccessToken(authNode.(*model.Node), payload) } else { result, err = service.RegisterNodeWithDiscovery(payload) } @@ -52,7 +52,7 @@ func AgentRegister(c *gin.Context) { // @Tags Agent // @Accept json // @Produce json -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Param payload body service.AgentNodePayload true "Agent heartbeat payload" // @Success 200 {object} map[string]interface{} // @Failure 400 {object} map[string]interface{} @@ -87,7 +87,7 @@ func AgentHeartbeat(c *gin.Context) { // @Tags Agent // @Accept json // @Produce json -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Param payload body service.AgentWAFIPGroupSyncInput true "WAF IP group sync payload" // @Success 200 {object} map[string]interface{} // @Failure 400 {object} map[string]interface{} @@ -109,7 +109,7 @@ func AgentSyncWAFIPGroups(c *gin.Context) { // @Summary Get active config for agent // @Tags Agent // @Produce json -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Success 200 {object} map[string]interface{} // @Router /api/agent/config-versions/active [get] func AgentGetActiveConfig(c *gin.Context) { @@ -143,7 +143,7 @@ func AgentGetActiveConfig(c *gin.Context) { // @Tags Agent // @Accept json // @Produce json -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Param payload body service.ApplyLogPayload true "Apply log payload" // @Success 200 {object} map[string]interface{} // @Failure 400 {object} map[string]interface{} @@ -169,7 +169,7 @@ func AgentReportApplyLog(c *gin.Context) { // AgentWebSocket godoc // @Summary Upgrade agent connection to websocket // @Tags Agent -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Router /api/agent/ws [get] func AgentWebSocket(c *gin.Context) { authNode, ok := c.Get("agent_node") diff --git a/openflare_server/controller/relay.go b/openflare_server/controller/relay.go index a83effcd..3c4fbc2a 100644 --- a/openflare_server/controller/relay.go +++ b/openflare_server/controller/relay.go @@ -16,7 +16,7 @@ import ( // @Tags Relay // @Accept json // @Produce json -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Param payload body service.RelayHeartbeatPayload true "Relay heartbeat payload" // @Success 200 {object} map[string]interface{} // @Failure 400 {object} map[string]interface{} @@ -44,7 +44,7 @@ func RelayHeartbeat(c *gin.Context) { // RelayWebSocket godoc // @Summary Upgrade relay connection to websocket // @Tags Relay -// @Security AgentTokenAuth +// @Security AccessTokenAuth // @Router /api/relay/ws [get] func RelayWebSocket(c *gin.Context) { authNode, ok := c.Get("relay_node") diff --git a/openflare_server/docs/docs.go b/openflare_server/docs/docs.go index e876e41b..1a66ccfb 100644 --- a/openflare_server/docs/docs.go +++ b/openflare_server/docs/docs.go @@ -355,7 +355,7 @@ const docTemplate = `{ "post": { "security": [ { - "AgentTokenAuth": [] + "AccessTokenAuth": [] } ], "consumes": [ @@ -401,7 +401,7 @@ const docTemplate = `{ "get": { "security": [ { - "AgentTokenAuth": [] + "AccessTokenAuth": [] } ], "produces": [ @@ -426,7 +426,7 @@ const docTemplate = `{ "post": { "security": [ { - "AgentTokenAuth": [] + "AccessTokenAuth": [] } ], "consumes": [ @@ -472,7 +472,7 @@ const docTemplate = `{ "post": { "security": [ { - "AgentTokenAuth": [] + "AccessTokenAuth": [] } ], "consumes": [ @@ -2801,7 +2801,7 @@ const docTemplate = `{ "$ref": "#/definitions/service.AgentNodeAccessLog" } }, - "agent_version": { + "version": { "type": "string" }, "buffered_observability": { @@ -2828,7 +2828,7 @@ const docTemplate = `{ "name": { "type": "string" }, - "nginx_version": { + "ext_version": { "type": "string" }, "node_id": { @@ -3186,7 +3186,7 @@ const docTemplate = `{ } }, "securityDefinitions": { - "AgentTokenAuth": { + "AccessTokenAuth": { "description": "Agent API 使用节点专属 Agent Token 或全局 Discovery Token", "type": "apiKey", "name": "X-Agent-Token", diff --git a/openflare_server/main.go b/openflare_server/main.go index 99f9d5d7..91c6e946 100644 --- a/openflare_server/main.go +++ b/openflare_server/main.go @@ -37,7 +37,7 @@ var indexPage []byte // @in header // @name Authorization // @description 管理端可使用 Bearer Token,例如:Bearer -// @securityDefinitions.apikey AgentTokenAuth +// @securityDefinitions.apikey AccessTokenAuth // @in header // @name X-Agent-Token // @description Agent API 使用节点专属 Agent Token 或全局 Discovery Token @@ -103,7 +103,7 @@ func main() { if common.SQLDSN != "" { dbBackend = "postgres" } - slog.Info("server config", "port", port, "gin_mode", gin.Mode(), "log_level", common.GetLogLevel(), "db_backend", dbBackend, "sqlite_path", common.SQLitePath, "redis_enabled", common.RedisEnabled, "log_dir", valueOrDefault(*common.LogDir, "stdout"), "agent_token_configured", common.AgentToken != "", "node_offline_threshold", common.NodeOfflineThreshold) + slog.Info("server config", "port", port, "gin_mode", gin.Mode(), "log_level", common.GetLogLevel(), "db_backend", dbBackend, "sqlite_path", common.SQLitePath, "redis_enabled", common.RedisEnabled, "log_dir", valueOrDefault(*common.LogDir, "stdout"), "access_token_configured", common.AccessToken != "", "node_offline_threshold", common.NodeOfflineThreshold) slog.Info("server listening", "address", fmt.Sprintf(":%s", port)) err = server.Run(":" + port) if err != nil { diff --git a/openflare_server/middleware/agent-auth.go b/openflare_server/middleware/agent-auth.go index 59ac3fce..368c1963 100644 --- a/openflare_server/middleware/agent-auth.go +++ b/openflare_server/middleware/agent-auth.go @@ -9,7 +9,7 @@ import ( func AgentAuth() func(c *gin.Context) { return func(c *gin.Context) { token := c.GetHeader("X-Agent-Token") - node, err := service.AuthenticateAgentToken(token) + node, err := service.AuthenticateAccessToken(token) if err != nil { c.JSON(http.StatusUnauthorized, gin.H{ "success": false, @@ -26,7 +26,7 @@ func AgentAuth() func(c *gin.Context) { func AgentRegisterAuth() func(c *gin.Context) { return func(c *gin.Context) { token := c.GetHeader("X-Agent-Token") - if node, err := service.AuthenticateAgentToken(token); err == nil { + if node, err := service.AuthenticateAccessToken(token); err == nil { c.Set("agent_node", node) c.Next() return diff --git a/openflare_server/middleware/relay-auth.go b/openflare_server/middleware/relay-auth.go index bdbc54a9..7e9d72c3 100644 --- a/openflare_server/middleware/relay-auth.go +++ b/openflare_server/middleware/relay-auth.go @@ -11,7 +11,7 @@ import ( func RelayAuth() func(c *gin.Context) { return func(c *gin.Context) { token := c.GetHeader("X-Agent-Token") - node, err := service.AuthenticateAgentToken(token) + node, err := service.AuthenticateAccessToken(token) if err != nil { c.JSON(http.StatusUnauthorized, gin.H{ "success": false, diff --git a/openflare_server/model/main.go b/openflare_server/model/main.go index e5e1f4cb..852796ea 100644 --- a/openflare_server/model/main.go +++ b/openflare_server/model/main.go @@ -41,6 +41,9 @@ func registeredModels() []any { &NodeRequestReport{}, &NodeAccessLog{}, &NodeHealthEvent{}, + &NodeObservationOpenresty{}, + &NodeObservationFrps{}, + &NodeObservationFrpc{}, &TLSCertificate{}, &ManagedDomain{}, &AcmeAccount{}, diff --git a/openflare_server/model/migrate/v20.go b/openflare_server/model/migrate/v20.go index 6e312619..40ce145a 100644 --- a/openflare_server/model/migrate/v20.go +++ b/openflare_server/model/migrate/v20.go @@ -32,6 +32,16 @@ func V20() Migration { } func migrateV20(ctx Context, db *gorm.DB, backend string) error { + if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_connections") { + if err := db.Migrator().AddColumn(&nodeV20{}, "RelayFrpsConnections"); err != nil { + return err + } + } + if !db.Migrator().HasColumn(&nodeV20{}, "relay_frps_proxy_count") { + if err := db.Migrator().AddColumn(&nodeV20{}, "RelayFrpsProxyCount"); err != nil { + return err + } + } if err := ctx.ApplyCurrentSchema(db, backend); err != nil { return err } diff --git a/openflare_server/model/migrate/v21.go b/openflare_server/model/migrate/v21.go new file mode 100644 index 00000000..95e15d81 --- /dev/null +++ b/openflare_server/model/migrate/v21.go @@ -0,0 +1,82 @@ +// v21 renames agent_token to access_token, unifies versions, and separates node observabilities. +package migrate + +import ( + "fmt" + "log/slog" + + "gorm.io/gorm" +) + +type nodeV21 struct{} + +func (nodeV21) TableName() string { + return "nodes" +} + +func init() { + Register(V21()) +} + +func V21() Migration { + return Migration{ + FromVersion: 20, + ToVersion: 21, + Migrate: migrateV21, + Validate: validateV21, + } +} + +func migrateV21(ctx Context, db *gorm.DB, backend string) error { + slog.Info("starting v21 database migration (Node Optimization & Observation Split)") + + migrator := db.Migrator() + + if migrator.HasColumn(&nodeV21{}, "agent_token") { + if err := migrator.RenameColumn(&nodeV21{}, "agent_token", "access_token"); err != nil { + return fmt.Errorf("failed to rename agent_token to access_token: %w", err) + } + } + + if migrator.HasColumn(&nodeV21{}, "agent_version") { + if err := migrator.RenameColumn(&nodeV21{}, "agent_version", "version"); err != nil { + return fmt.Errorf("failed to rename agent_version to version: %w", err) + } + } + + if migrator.HasColumn(&nodeV21{}, "nginx_version") { + if err := migrator.RenameColumn(&nodeV21{}, "nginx_version", "ext_version"); err != nil { + return fmt.Errorf("failed to rename nginx_version to ext_version: %w", err) + } + } + + // Drop old merged columns + columnsToDrop := []string{ + "relay_version", + "relay_frp_version", + "relay_frps_connections", + "relay_frps_proxy_count", + } + + for _, col := range columnsToDrop { + if migrator.HasColumn(&nodeV21{}, col) { + if err := migrator.DropColumn(&nodeV21{}, col); err != nil { + slog.Warn("failed to drop column in v21 migration", "column", col, "error", err) + } + } + } + + if err := ctx.ApplyCurrentSchema(db, backend); err != nil { + return err + } + + slog.Info("completed v21 database migration") + return validateV21(ctx, db, backend) +} + +func validateV21(ctx Context, db *gorm.DB, backend string) error { + if err := ctx.ValidateDatabaseSchemaVersion(db, backend, 21); err != nil { + return err + } + return nil +} diff --git a/openflare_server/model/migrations.go b/openflare_server/model/migrations.go index 11645087..8c737858 100644 --- a/openflare_server/model/migrations.go +++ b/openflare_server/model/migrations.go @@ -86,6 +86,8 @@ func (databaseSchemaMigrationContext) ValidateDatabaseSchemaVersion(db *gorm.DB, return validateDatabaseSchemaV19(db, backend) case 20: return validateDatabaseSchemaV20(db, backend) + case 21: + return validateDatabaseSchemaV21(db, backend) default: return fmt.Errorf("database schema validation for v%d is not defined", version) } @@ -1264,11 +1266,8 @@ func validateExternalDatabaseSchema(ctx databaseSchemaMigrationContext, db *gorm return ctx.ValidateDatabaseSchemaVersion(db, backend, targetVersion) } for _, migration := range schemamigrate.Migrations() { - if migration.ToVersion > targetVersion { - continue - } - if err := migration.Validate(ctx, db, backend); err != nil { - return err + if migration.ToVersion == targetVersion { + return migration.Validate(ctx, db, backend) } } return nil @@ -1373,6 +1372,22 @@ func initializeFreshDatabaseSchema(db *gorm.DB, backend string) error { return saveDatabaseSchemaVersion(db, currentDatabaseSchemaVersion) } +func validateDatabaseSchemaV21(db *gorm.DB, backend string) error { + if err := validateDatabaseSchemaV19(db, backend); err != nil { + return err + } + if !db.Migrator().HasColumn(&Node{}, "access_token") { + return fmt.Errorf("column nodes.access_token is missing") + } + if !db.Migrator().HasColumn(&Node{}, "version") { + return fmt.Errorf("column nodes.version is missing") + } + if !db.Migrator().HasColumn(&Node{}, "ext_version") { + return fmt.Errorf("column nodes.ext_version is missing") + } + return nil +} + func ensureDatabaseSchemaUpToDate(db *gorm.DB, backend string) error { version, exists, err := loadDatabaseSchemaVersion(db) if err != nil { diff --git a/openflare_server/model/node.go b/openflare_server/model/node.go index e275a708..84222358 100644 --- a/openflare_server/model/node.go +++ b/openflare_server/model/node.go @@ -12,14 +12,14 @@ type Node struct { GeoLatitude *float64 `json:"geo_latitude"` GeoLongitude *float64 `json:"geo_longitude"` GeoManualOverride bool `json:"geo_manual_override" gorm:"not null;default:false"` - AgentToken string `json:"-" gorm:"size:128;index"` + AccessToken string `json:"-" gorm:"column:access_token;size:128;index"` AutoUpdateEnabled bool `json:"auto_update_enabled" gorm:"not null;default:false"` UpdateRequested bool `json:"update_requested" gorm:"not null;default:false"` UpdateChannel string `json:"update_channel" gorm:"size:16;not null;default:'stable'"` UpdateTag string `json:"update_tag" gorm:"size:64"` RestartOpenrestyRequested bool `json:"restart_openresty_requested" gorm:"not null;default:false"` - AgentVersion string `json:"agent_version" gorm:"size:64;not null"` - NginxVersion string `json:"nginx_version" gorm:"size:64"` + Version string `json:"version" gorm:"size:64;not null"` + ExtVersion string `json:"ext_version" gorm:"size:64"` OpenrestyStatus string `json:"openresty_status" gorm:"size:16;not null;default:'unknown'"` OpenrestyMessage string `json:"openresty_message" gorm:"type:text"` Status string `json:"status" gorm:"size:16;not null;default:'offline'"` @@ -28,7 +28,7 @@ type Node struct { LastError string `json:"last_error" gorm:"type:text"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` - // Node type: edge_node (default) | tunnel_relay + // Node type: edge_node (default) | tunnel_relay | tunnel_client NodeType string `json:"node_type" gorm:"size:32;not null;default:'edge_node'"` // TunnelRelay specific fields RelayBindPort int `json:"relay_bind_port" gorm:"not null;default:0"` @@ -38,10 +38,6 @@ type Node struct { RelayClientAccessAddr string `json:"relay_client_access_addr" gorm:"size:255"` RelayClientProxyURL string `json:"relay_client_proxy_url" gorm:"size:512"` 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) { @@ -69,9 +65,9 @@ func GetNodeByID(id uint) (*Node, error) { return node, err } -func GetNodeByAgentToken(token string) (*Node, error) { +func GetNodeByAccessToken(token string) (*Node, error) { node := &Node{} - err := DB.Where("agent_token = ?", token).First(node).Error + err := DB.Where("access_token = ?", token).First(node).Error return node, err } diff --git a/openflare_server/model/node_metric_snapshot.go b/openflare_server/model/node_metric_snapshot.go index e71d4bcc..66764a96 100644 --- a/openflare_server/model/node_metric_snapshot.go +++ b/openflare_server/model/node_metric_snapshot.go @@ -8,22 +8,19 @@ import ( ) type NodeMetricSnapshot struct { - ID uint `json:"id" gorm:"primaryKey"` - NodeID string `json:"node_id" gorm:"index;size:64;not null"` - CapturedAt time.Time `json:"captured_at" gorm:"index"` - CPUUsagePercent float64 `json:"cpu_usage_percent"` - MemoryUsedBytes int64 `json:"memory_used_bytes"` - MemoryTotalBytes int64 `json:"memory_total_bytes"` - StorageUsedBytes int64 `json:"storage_used_bytes"` - StorageTotalBytes int64 `json:"storage_total_bytes"` - DiskReadBytes int64 `json:"disk_read_bytes"` - DiskWriteBytes int64 `json:"disk_write_bytes"` - NetworkRxBytes int64 `json:"network_rx_bytes"` - NetworkTxBytes int64 `json:"network_tx_bytes"` - OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` - OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` - OpenrestyConnections int64 `json:"openresty_connections"` - CreatedAt time.Time `json:"created_at"` + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"index;size:64;not null"` + CapturedAt time.Time `json:"captured_at" gorm:"index"` + CPUUsagePercent float64 `json:"cpu_usage_percent"` + MemoryUsedBytes int64 `json:"memory_used_bytes"` + MemoryTotalBytes int64 `json:"memory_total_bytes"` + StorageUsedBytes int64 `json:"storage_used_bytes"` + StorageTotalBytes int64 `json:"storage_total_bytes"` + DiskReadBytes int64 `json:"disk_read_bytes"` + DiskWriteBytes int64 `json:"disk_write_bytes"` + NetworkRxBytes int64 `json:"network_rx_bytes"` + NetworkTxBytes int64 `json:"network_tx_bytes"` + CreatedAt time.Time `json:"created_at"` } func (snapshot *NodeMetricSnapshot) GetID() uint { diff --git a/openflare_server/model/node_observation_frpc.go b/openflare_server/model/node_observation_frpc.go new file mode 100644 index 00000000..c8beb2e4 --- /dev/null +++ b/openflare_server/model/node_observation_frpc.go @@ -0,0 +1,60 @@ +package model + +import ( + "openflare/utils" + "time" + + "gorm.io/gorm" +) + +type NodeObservationFrpc struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"index;size:64;not null"` + CapturedAt time.Time `json:"captured_at" gorm:"index"` + TunnelStatus string `json:"tunnel_status" gorm:"size:16"` + ConnectedRelaysCount int `json:"connected_relays_count"` + CreatedAt time.Time `json:"created_at"` +} + +func (obs *NodeObservationFrpc) GetID() uint { + return obs.ID +} + +func (obs *NodeObservationFrpc) GetTime() time.Time { + return obs.CapturedAt +} + +func (obs *NodeObservationFrpc) BeforeCreate(tx *gorm.DB) error { + return assignObservabilityID(&obs.ID) +} + +func (obs *NodeObservationFrpc) Insert() error { + return DB.Create(obs).Error +} + +func ListNodeObservationFrpcs(nodeID string, since time.Time, limit int) (observations []*NodeObservationFrpc, err error) { + rows, err := queryAcrossShards("node_observation_frpcs", func(tx *gorm.DB) ([]*NodeObservationFrpc, error) { + var shardRows []*NodeObservationFrpc + query := tx.Order("captured_at desc, id desc") + if nodeID != "" { + query = query.Where("node_id = ?", nodeID) + } + if !since.IsZero() { + query = query.Where("captured_at >= ?", since) + } + if err := query.Find(&shardRows).Error; err != nil { + return nil, err + } + return shardRows, nil + }) + if err != nil { + return nil, err + } + return utils.SortAndLimitRecords(rows, limit), nil +} + +func DeleteNodeObservationFrpcsBefore(db *gorm.DB, before time.Time) (int64, error) { + return deleteAcrossShards(db, "node_observation_frpcs", &NodeObservationFrpc{}, func(tx *gorm.DB) *gorm.DB { + return tx.Where("captured_at < ?", before) + }) +} diff --git a/openflare_server/model/node_observation_frps.go b/openflare_server/model/node_observation_frps.go new file mode 100644 index 00000000..5e80bf0e --- /dev/null +++ b/openflare_server/model/node_observation_frps.go @@ -0,0 +1,60 @@ +package model + +import ( + "openflare/utils" + "time" + + "gorm.io/gorm" +) + +type NodeObservationFrps struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"index;size:64;not null"` + CapturedAt time.Time `json:"captured_at" gorm:"index"` + FrpsConnections int `json:"frps_connections"` + FrpsProxyCount int `json:"frps_proxy_count"` + CreatedAt time.Time `json:"created_at"` +} + +func (obs *NodeObservationFrps) GetID() uint { + return obs.ID +} + +func (obs *NodeObservationFrps) GetTime() time.Time { + return obs.CapturedAt +} + +func (obs *NodeObservationFrps) BeforeCreate(tx *gorm.DB) error { + return assignObservabilityID(&obs.ID) +} + +func (obs *NodeObservationFrps) Insert() error { + return DB.Create(obs).Error +} + +func ListNodeObservationFrps(nodeID string, since time.Time, limit int) (observations []*NodeObservationFrps, err error) { + rows, err := queryAcrossShards("node_observation_frps", func(tx *gorm.DB) ([]*NodeObservationFrps, error) { + var shardRows []*NodeObservationFrps + query := tx.Order("captured_at desc, id desc") + if nodeID != "" { + query = query.Where("node_id = ?", nodeID) + } + if !since.IsZero() { + query = query.Where("captured_at >= ?", since) + } + if err := query.Find(&shardRows).Error; err != nil { + return nil, err + } + return shardRows, nil + }) + if err != nil { + return nil, err + } + return utils.SortAndLimitRecords(rows, limit), nil +} + +func DeleteNodeObservationFrpsBefore(db *gorm.DB, before time.Time) (int64, error) { + return deleteAcrossShards(db, "node_observation_frps", &NodeObservationFrps{}, func(tx *gorm.DB) *gorm.DB { + return tx.Where("captured_at < ?", before) + }) +} diff --git a/openflare_server/model/node_observation_openresty.go b/openflare_server/model/node_observation_openresty.go new file mode 100644 index 00000000..2a14bf3d --- /dev/null +++ b/openflare_server/model/node_observation_openresty.go @@ -0,0 +1,61 @@ +package model + +import ( + "openflare/utils" + "time" + + "gorm.io/gorm" +) + +type NodeObservationOpenresty struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"index;size:64;not null"` + CapturedAt time.Time `json:"captured_at" gorm:"index"` + OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` + OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` + OpenrestyConnections int64 `json:"openresty_connections"` + CreatedAt time.Time `json:"created_at"` +} + +func (obs *NodeObservationOpenresty) GetID() uint { + return obs.ID +} + +func (obs *NodeObservationOpenresty) GetTime() time.Time { + return obs.CapturedAt +} + +func (obs *NodeObservationOpenresty) BeforeCreate(tx *gorm.DB) error { + return assignObservabilityID(&obs.ID) +} + +func (obs *NodeObservationOpenresty) Insert() error { + return DB.Create(obs).Error +} + +func ListNodeObservationOpenresty(nodeID string, since time.Time, limit int) (observations []*NodeObservationOpenresty, err error) { + rows, err := queryAcrossShards("node_observation_openresties", func(tx *gorm.DB) ([]*NodeObservationOpenresty, error) { + var shardRows []*NodeObservationOpenresty + query := tx.Order("captured_at desc, id desc") + if nodeID != "" { + query = query.Where("node_id = ?", nodeID) + } + if !since.IsZero() { + query = query.Where("captured_at >= ?", since) + } + if err := query.Find(&shardRows).Error; err != nil { + return nil, err + } + return shardRows, nil + }) + if err != nil { + return nil, err + } + return utils.SortAndLimitRecords(rows, limit), nil +} + +func DeleteNodeObservationOpenrestiesBefore(db *gorm.DB, before time.Time) (int64, error) { + return deleteAcrossShards(db, "node_observation_openresties", &NodeObservationOpenresty{}, func(tx *gorm.DB) *gorm.DB { + return tx.Where("captured_at < ?", before) + }) +} diff --git a/openflare_server/model/sharding.go b/openflare_server/model/sharding.go index d5d64791..a615c7bc 100644 --- a/openflare_server/model/sharding.go +++ b/openflare_server/model/sharding.go @@ -48,6 +48,9 @@ func shardedObservabilityTables() []any { &NodeMetricSnapshot{}, &NodeRequestReport{}, &NodeAccessLog{}, + &NodeObservationOpenresty{}, + &NodeObservationFrps{}, + &NodeObservationFrpc{}, } } @@ -56,12 +59,15 @@ func shardedObservabilityBaseTables() []string { "node_metric_snapshots", "node_request_reports", "node_access_logs", + "node_observation_openresties", + "node_observation_frps", + "node_observation_frpcs", } } func isShardedObservabilityTable(tableName string) bool { switch strings.TrimSpace(tableName) { - case "node_metric_snapshots", "node_request_reports", "node_access_logs": + case "node_metric_snapshots", "node_request_reports", "node_access_logs", "node_observation_openresties", "node_observation_frps", "node_observation_frpcs": return true default: return false diff --git a/openflare_server/router/api_phase1_test.go b/openflare_server/router/api_phase1_test.go index e1275af1..cf979daf 100644 --- a/openflare_server/router/api_phase1_test.go +++ b/openflare_server/router/api_phase1_test.go @@ -363,19 +363,19 @@ func TestPhase1HTTPSAndCertificateImportLifecycle(t *testing.T) { t.Fatal("expected support files json to contain certificate artifacts") } if err := (&model.Node{ - NodeID: "phase1-node", - Name: "phase1-node", - IP: "10.0.0.8", - AgentToken: common.AgentToken, - AgentVersion: "0.1.0", - NginxVersion: "1.25.5", - Status: service.NodeStatusOnline, - LastSeenAt: time.Now(), + NodeID: "phase1-node", + Name: "phase1-node", + IP: "10.0.0.8", + AccessToken: common.AccessToken, + Version: "0.1.0", + ExtVersion: "1.25.5", + Status: service.NodeStatusOnline, + LastSeenAt: time.Now(), }).Insert(); err != nil { t.Fatalf("failed to seed phase1 node: %v", err) } - agentResp := performAgentJSONRequestWithToken(t, engine, common.AgentToken, http.MethodGet, "/api/agent/config-versions/active", nil) + agentResp := performAgentJSONRequestWithToken(t, engine, common.AccessToken, http.MethodGet, "/api/agent/config-versions/active", nil) var activeConfig map[string]any decodeResponseData(t, agentResp, &activeConfig) sourceConfigJSON, ok := activeConfig["source_config_json"].(string) @@ -481,7 +481,7 @@ func setupTestDB(t *testing.T) { t.Helper() dbPath := filepath.Join(t.TempDir(), "phase1.db") common.SQLitePath = dbPath - common.AgentToken = "phase1-agent-token" + common.AccessToken = "phase1-agent-token" if err := model.InitDB(); err != nil { t.Fatalf("failed to init db: %v", err) } diff --git a/openflare_server/router/api_phase2_test.go b/openflare_server/router/api_phase2_test.go index 2359d5d6..4ddbe147 100644 --- a/openflare_server/router/api_phase2_test.go +++ b/openflare_server/router/api_phase2_test.go @@ -426,7 +426,7 @@ func TestPhase2AgentLifecycle(t *testing.T) { }) var createdNode service.NodeView decodeResponseData(t, createdNodeResp, &createdNode) - if createdNode.AgentToken == "" || createdNode.Status != service.NodeStatusPending { + if createdNode.AccessToken == "" || createdNode.Status != service.NodeStatusPending { t.Fatal("expected created node to expose agent token with pending status") } if createdNode.GeoName != "Shanghai" || createdNode.GeoLatitude == nil || createdNode.GeoLongitude == nil { @@ -437,31 +437,31 @@ func TestPhase2AgentLifecycle(t *testing.T) { "node_id": "spoofed-node-id", "name": "shanghai-edge-1", "ip": "10.0.0.9", - "agent_version": "0.1.1", - "nginx_version": "1.27.1.2", + "version": "0.1.1", + "ext_version": "1.27.1.2", "openresty_status": service.OpenrestyStatusUnhealthy, "openresty_message": "docker run openresty failed: bind 80 already allocated", "current_version": "", "last_error": "", } - resp := performAgentJSONRequestWithTokenAndRemote(t, engine, createdNode.AgentToken, http.MethodPost, "/api/agent/nodes/heartbeat", heartbeatPayload, "198.51.100.10:1234") + resp := performAgentJSONRequestWithTokenAndRemote(t, engine, createdNode.AccessToken, http.MethodPost, "/api/agent/nodes/heartbeat", heartbeatPayload, "198.51.100.10:1234") var registeredNode model.Node decodeResponseData(t, resp, ®isteredNode) - if registeredNode.IP != "198.51.100.10" || registeredNode.AgentVersion != "0.1.1" || registeredNode.NodeID != createdNode.NodeID { + if registeredNode.IP != "198.51.100.10" || registeredNode.Version != "0.1.1" || registeredNode.NodeID != createdNode.NodeID { t.Fatal("expected heartbeat to update node metadata") } if registeredNode.OpenrestyStatus != service.OpenrestyStatusUnhealthy { t.Fatal("expected heartbeat to update openresty status") } - activeConfigResp := performAgentJSONRequestWithToken(t, engine, createdNode.AgentToken, http.MethodGet, "/api/agent/config-versions/active", nil) + activeConfigResp := performAgentJSONRequestWithToken(t, engine, createdNode.AccessToken, http.MethodGet, "/api/agent/config-versions/active", nil) var activeConfig service.AgentConfigResponse decodeResponseData(t, activeConfigResp, &activeConfig) if activeConfig.Version == "" || activeConfig.SourceConfigJSON == "" || activeConfig.Checksum == "" { t.Fatal("expected active config response to contain version payload") } - successApplyResp := performAgentJSONRequestWithToken(t, engine, createdNode.AgentToken, http.MethodPost, "/api/agent/apply-logs", map[string]any{ + successApplyResp := performAgentJSONRequestWithToken(t, engine, createdNode.AccessToken, http.MethodPost, "/api/agent/apply-logs", map[string]any{ "node_id": "spoofed-node-id", "version": activeConfig.Version, "result": service.ApplyResultOK, @@ -473,7 +473,7 @@ func TestPhase2AgentLifecycle(t *testing.T) { t.Fatal("expected apply log success to be recorded") } - failedApplyResp := performAgentJSONRequestWithToken(t, engine, createdNode.AgentToken, http.MethodPost, "/api/agent/apply-logs", map[string]any{ + failedApplyResp := performAgentJSONRequestWithToken(t, engine, createdNode.AccessToken, http.MethodPost, "/api/agent/apply-logs", map[string]any{ "node_id": "spoofed-node-id", "version": activeConfig.Version, "result": service.ApplyResultFailed, @@ -494,7 +494,7 @@ func TestPhase2AgentLifecycle(t *testing.T) { if nodes[0].Status != service.NodeStatusOnline { t.Fatal("expected registered node to become online") } - if nodes[0].AgentToken != createdNode.AgentToken { + if nodes[0].AccessToken != createdNode.AccessToken { t.Fatal("expected node auth token to remain stable after occupancy") } if nodes[0].LatestApplyResult != service.ApplyResultFailed || nodes[0].LatestApplyMessage != "openresty reload failed" { @@ -561,7 +561,7 @@ func TestPhase2AgentLifecycle(t *testing.T) { } restartHeartbeatReq := httptest.NewRequest(http.MethodPost, "/api/agent/nodes/heartbeat", bytes.NewReader(rawHeartbeatPayload)) restartHeartbeatReq.Header.Set("Content-Type", "application/json") - restartHeartbeatReq.Header.Set("X-Agent-Token", createdNode.AgentToken) + restartHeartbeatReq.Header.Set("X-Agent-Token", createdNode.AccessToken) restartHeartbeatReq.RemoteAddr = "198.51.100.10:1234" restartHeartbeatRecorder := httptest.NewRecorder() engine.ServeHTTP(restartHeartbeatRecorder, restartHeartbeatReq) @@ -631,7 +631,7 @@ func TestPhase2AgentLifecycle(t *testing.T) { if logs.Total != 0 || len(logs.Rows) != 0 || logs.Current != 1 || logs.TotalPage != 0 { t.Fatalf("expected empty apply log page after delete-all cleanup, got %+v", logs) } - postDeleteApplyResp := performAgentJSONRequestWithToken(t, engine, createdNode.AgentToken, http.MethodPost, "/api/agent/apply-logs", map[string]any{ + postDeleteApplyResp := performAgentJSONRequestWithToken(t, engine, createdNode.AccessToken, http.MethodPost, "/api/agent/apply-logs", map[string]any{ "version": activeConfig.Version, "result": service.ApplyResultOK, "message": "local config already matches active version; apply skipped", @@ -678,9 +678,9 @@ func TestPhase2AgentLifecycle(t *testing.T) { t.Fatalf("expected delete node success, got %s", deleteResp.Message) } - deniedReq := httptest.NewRequest(http.MethodPost, "/api/agent/nodes/heartbeat", bytes.NewReader([]byte(`{"ip":"10.0.0.9","agent_version":"0.1.1"}`))) + deniedReq := httptest.NewRequest(http.MethodPost, "/api/agent/nodes/heartbeat", bytes.NewReader([]byte(`{"ip":"10.0.0.9","version":"0.1.1"}`))) deniedReq.Header.Set("Content-Type", "application/json") - deniedReq.Header.Set("X-Agent-Token", createdNode.AgentToken) + deniedReq.Header.Set("X-Agent-Token", createdNode.AccessToken) deniedRecorder := httptest.NewRecorder() engine.ServeHTTP(deniedRecorder, deniedReq) if deniedRecorder.Code != http.StatusUnauthorized { @@ -853,14 +853,14 @@ func TestPhase2GlobalDiscoveryRegistration(t *testing.T) { "node_id": "local-node-id", "name": "bulk-edge-1", "ip": "10.0.0.18", - "agent_version": "0.2.0", - "nginx_version": "1.25.5", + "version": "0.2.0", + "ext_version": "1.25.5", "current_version": "", "last_error": "", }, "203.0.113.18:4321") var registration service.AgentRegistrationResponse decodeResponseData(t, resp, ®istration) - if registration.AgentToken == "" || registration.NodeID == "" { + if registration.AccessToken == "" || registration.NodeID == "" { t.Fatal("expected discovery registration to issue node-specific agent token") } @@ -870,7 +870,7 @@ func TestPhase2GlobalDiscoveryRegistration(t *testing.T) { if len(nodes) != 1 { t.Fatalf("expected 1 discovered node, got %d", len(nodes)) } - if nodes[0].Name != "bulk-edge-1" || nodes[0].AgentToken != registration.AgentToken || nodes[0].Status != service.NodeStatusOnline { + if nodes[0].Name != "bulk-edge-1" || nodes[0].AccessToken != registration.AccessToken || nodes[0].Status != service.NodeStatusOnline { t.Fatal("expected discovered node to be created online with issued agent token") } if nodes[0].IP != "203.0.113.18" { diff --git a/openflare_server/service/agent.go b/openflare_server/service/agent.go index 8b355a95..a487e891 100644 --- a/openflare_server/service/agent.go +++ b/openflare_server/service/agent.go @@ -29,14 +29,15 @@ type AgentNodePayload struct { NodeID string `json:"node_id"` Name string `json:"name"` IP string `json:"ip"` - AgentVersion string `json:"agent_version"` - NginxVersion string `json:"nginx_version"` + Version string `json:"version"` + ExtVersion string `json:"ext_version"` CurrentVersion string `json:"current_version"` LastError string `json:"last_error"` OpenrestyStatus string `json:"openresty_status"` OpenrestyMessage string `json:"openresty_message"` Profile *AgentNodeSystemProfile `json:"profile,omitempty"` Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"` + OpenrestyObservation *AgentNodeOpenrestyObservation `json:"openresty_observation,omitempty"` TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"` AccessLogs []AgentNodeAccessLog `json:"access_logs,omitempty"` BufferedObservability []AgentBufferedObservabilityRecord `json:"buffered_observability,omitempty"` @@ -139,14 +140,14 @@ type NodeView struct { GeoLatitude *float64 `json:"geo_latitude"` GeoLongitude *float64 `json:"geo_longitude"` GeoManualOverride bool `json:"geo_manual_override"` - AgentToken string `json:"agent_token"` + AccessToken string `json:"access_token"` AutoUpdateEnabled bool `json:"auto_update_enabled"` UpdateRequested bool `json:"update_requested"` UpdateChannel string `json:"update_channel"` UpdateTag string `json:"update_tag"` RestartOpenrestyRequested bool `json:"restart_openresty_requested"` - AgentVersion string `json:"agent_version"` - NginxVersion string `json:"nginx_version"` + Version string `json:"version"` + ExtVersion string `json:"ext_version"` OpenrestyStatus string `json:"openresty_status"` OpenrestyMessage string `json:"openresty_message"` Status string `json:"status"` @@ -170,10 +171,6 @@ type NodeView struct { RelayClientAccessAddr string `json:"relay_client_access_addr"` RelayClientProxyURL string `json:"relay_client_proxy_url"` 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) { @@ -199,7 +196,7 @@ func HeartbeatNode(node *model.Node, payload AgentNodePayload) (*HeartbeatRespon return nil, err } } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) persistHeartbeatObservability(node.NodeID, payload, node.LastSeenAt) activeConfig, err := GetActiveConfigMetaForAgent() if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { @@ -493,8 +490,8 @@ func collectNodeHeartbeatChanges(previous *model.Node, current *model.Node) map[ appendIfChanged("name", previous.Name, current.Name) appendIfChanged("ip", previous.IP, current.IP) appendIfChanged("geo_name", previous.GeoName, current.GeoName) - appendIfChanged("agent_version", previous.AgentVersion, current.AgentVersion) - appendIfChanged("nginx_version", previous.NginxVersion, current.NginxVersion) + appendIfChanged("version", previous.Version, current.Version) + appendIfChanged("ext_version", previous.ExtVersion, current.ExtVersion) appendIfChanged("openresty_status", previous.OpenrestyStatus, current.OpenrestyStatus) appendIfChanged("openresty_message", previous.OpenrestyMessage, current.OpenrestyMessage) appendIfChanged("status", previous.Status, current.Status) diff --git a/openflare_server/service/agent_test.go b/openflare_server/service/agent_test.go index c755e6c1..7a22fe89 100644 --- a/openflare_server/service/agent_test.go +++ b/openflare_server/service/agent_test.go @@ -177,7 +177,7 @@ func TestGetActiveConfigForAgentUsesTenMinutePoWSessionDefault(t *testing.T) { } } -func TestRegisterNodeWithAgentToken(t *testing.T) { +func TestRegisterNodeWithAccessToken(t *testing.T) { setupServiceTestDB(t) // 1. Success path @@ -203,17 +203,17 @@ func TestRegisterNodeWithAgentToken(t *testing.T) { payload := AgentNodePayload{ Name: "payload-name-should-be-ignored", IP: "192.168.1.20", - AgentVersion: "v1.0.1", - NginxVersion: "1.27.1.3", + Version: "v1.0.1", + ExtVersion: "1.27.1.3", OpenrestyStatus: "healthy", } - resp, err := RegisterNodeWithAgentToken(stored, payload) + resp, err := RegisterNodeWithAccessToken(stored, payload) if err != nil { - t.Fatalf("RegisterNodeWithAgentToken failed: %v", err) + t.Fatalf("RegisterNodeWithAccessToken failed: %v", err) } - if resp.NodeID != stored.NodeID || resp.AgentToken != stored.AgentToken || resp.Name != "reserved-node-1" { + if resp.NodeID != stored.NodeID || resp.AccessToken != stored.AccessToken || resp.Name != "reserved-node-1" { t.Errorf("unexpected response: %+v", resp) } @@ -222,7 +222,7 @@ func TestRegisterNodeWithAgentToken(t *testing.T) { if err != nil { t.Fatalf("failed to fetch updated node: %v", err) } - if updated.AgentVersion != "v1.0.1" || updated.NginxVersion != "1.27.1.3" || updated.OpenrestyStatus != "healthy" { + if updated.Version != "v1.0.1" || updated.ExtVersion != "1.27.1.3" || updated.OpenrestyStatus != "healthy" { t.Errorf("node attributes were not updated: %+v", updated) } // Name should be preserved since preserveName is true @@ -231,7 +231,7 @@ func TestRegisterNodeWithAgentToken(t *testing.T) { } // 2. Fail path - Nil Node - _, err = RegisterNodeWithAgentToken(nil, payload) + _, err = RegisterNodeWithAccessToken(nil, payload) if err == nil || !strings.Contains(err.Error(), "节点不存在") { t.Errorf("expected error '节点不存在', got %v", err) } @@ -239,16 +239,16 @@ func TestRegisterNodeWithAgentToken(t *testing.T) { // 3. Fail path - Invalid Payload (empty IP) badPayload := payload badPayload.IP = "" - _, err = RegisterNodeWithAgentToken(stored, badPayload) + _, err = RegisterNodeWithAccessToken(stored, badPayload) if err == nil || !strings.Contains(err.Error(), "ip 不能为空") { t.Errorf("expected error 'ip 不能为空', got %v", err) } // 4. Name update if empty emptyNameNode := &model.Node{ - NodeID: "node-empty-name", - Name: "", - AgentToken: "empty-name-token", + NodeID: "node-empty-name", + Name: "", + AccessToken: "empty-name-token", } if err := emptyNameNode.Insert(); err != nil { t.Fatalf("failed to insert emptyNameNode: %v", err) @@ -256,9 +256,9 @@ func TestRegisterNodeWithAgentToken(t *testing.T) { payloadWithName := payload payloadWithName.Name = "filled-name" payloadWithName.IP = "192.168.1.30" - _, err = RegisterNodeWithAgentToken(emptyNameNode, payloadWithName) + _, err = RegisterNodeWithAccessToken(emptyNameNode, payloadWithName) if err != nil { - t.Fatalf("RegisterNodeWithAgentToken empty name node failed: %v", err) + t.Fatalf("RegisterNodeWithAccessToken empty name node failed: %v", err) } updatedEmptyName, err := model.GetNodeByNodeID("node-empty-name") if err != nil { @@ -276,8 +276,8 @@ func TestRegisterNodeWithDiscovery(t *testing.T) { payload := AgentNodePayload{ Name: "discovery-node", IP: "192.168.2.10", - AgentVersion: "v1.0.0", - NginxVersion: "1.27.1.3", + Version: "v1.0.0", + ExtVersion: "1.27.1.3", OpenrestyStatus: "healthy", } @@ -286,7 +286,7 @@ func TestRegisterNodeWithDiscovery(t *testing.T) { t.Fatalf("RegisterNodeWithDiscovery failed: %v", err) } - if resp.NodeID == "" || resp.AgentToken == "" || resp.Name != "discovery-node" { + if resp.NodeID == "" || resp.AccessToken == "" || resp.Name != "discovery-node" { t.Errorf("unexpected response: %+v", resp) } @@ -295,7 +295,7 @@ func TestRegisterNodeWithDiscovery(t *testing.T) { if err != nil { t.Fatalf("failed to fetch node: %v", err) } - if node.IP != "192.168.2.10" || node.AgentVersion != "v1.0.0" || node.Name != "discovery-node" { + if node.IP != "192.168.2.10" || node.Version != "v1.0.0" || node.Name != "discovery-node" { t.Errorf("unexpected stored node data: %+v", node) } @@ -317,10 +317,10 @@ func TestRegisterNodeWithDiscovery(t *testing.T) { // 3. Fail path - Invalid Payload (empty AgentVersion) badPayload := payload - badPayload.AgentVersion = "" + badPayload.Version = "" _, err = RegisterNodeWithDiscovery(badPayload) - if err == nil || !strings.Contains(err.Error(), "agent_version 不能为空") { - t.Errorf("expected error 'agent_version 不能为空', got %v", err) + if err == nil || !strings.Contains(err.Error(), "version 不能为空") { + t.Errorf("expected error 'version 不能为空', got %v", err) } } @@ -329,12 +329,12 @@ func TestReportApplyLog_Success(t *testing.T) { // Seed node node := &model.Node{ - NodeID: "node-apply-1", - Name: "apply-edge", - IP: "192.168.3.10", - AgentToken: "apply-token", - AgentVersion: "v1.0.0", - Status: NodeStatusOffline, + NodeID: "node-apply-1", + Name: "apply-edge", + IP: "192.168.3.10", + AccessToken: "apply-token", + Version: "v1.0.0", + Status: NodeStatusOffline, } if err := node.Insert(); err != nil { t.Fatalf("failed to insert node: %v", err) @@ -387,8 +387,8 @@ func TestReportApplyLog_WarningAndFailure(t *testing.T) { NodeID: "node-apply-2", Name: "apply-edge-2", IP: "192.168.3.20", - AgentToken: "apply-token-2", - AgentVersion: "v1.0.0", + AccessToken: "apply-token-2", + Version: "v1.0.0", CurrentVersion: "20260531-001", // Old version Status: NodeStatusOnline, } @@ -452,12 +452,12 @@ func TestReportApplyLog_Failures(t *testing.T) { // Seed node node := &model.Node{ - NodeID: "node-apply-3", - Name: "apply-edge-3", - IP: "192.168.3.30", - AgentToken: "apply-token-3", - AgentVersion: "v1.0.0", - Status: NodeStatusOnline, + NodeID: "node-apply-3", + Name: "apply-edge-3", + IP: "192.168.3.30", + AccessToken: "apply-token-3", + Version: "v1.0.0", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to insert node: %v", err) @@ -508,11 +508,11 @@ func TestListAndCleanupApplyLogs(t *testing.T) { // Seed node node := &model.Node{ - NodeID: "node-logs", - Name: "logs-edge", - IP: "192.168.4.10", - AgentToken: "logs-token", - Status: NodeStatusOnline, + NodeID: "node-logs", + Name: "logs-edge", + IP: "192.168.4.10", + AccessToken: "logs-token", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to insert node: %v", err) diff --git a/openflare_server/service/dashboard.go b/openflare_server/service/dashboard.go index e4155033..45c6595b 100644 --- a/openflare_server/service/dashboard.go +++ b/openflare_server/service/dashboard.go @@ -93,6 +93,11 @@ func GetDashboardOverview() (*DashboardOverviewView, error) { return nil, err } + openrestySnapshots, err := model.ListNodeObservationOpenresty("", since, 0) + if err != nil { + return nil, err + } + view := &DashboardOverviewView{ GeneratedAt: now, Nodes: make([]DashboardNodeHealth, 0, len(nodes)), @@ -100,7 +105,7 @@ func GetDashboardOverview() (*DashboardOverviewView, error) { Trends: DashboardTrends{ Traffic24h: buildTrafficTrendPoints(now, reports), Capacity24h: buildCapacityTrendPoints(now, snapshots), - Network24h: buildNetworkTrendPoints(now, snapshots), + Network24h: buildNetworkTrendPoints(now, snapshots, openrestySnapshots), DiskIO24h: buildDiskIOTrendPoints(now, snapshots), }, } diff --git a/openflare_server/service/https_phase1_test.go b/openflare_server/service/https_phase1_test.go index a59007b1..1fe08fce 100644 --- a/openflare_server/service/https_phase1_test.go +++ b/openflare_server/service/https_phase1_test.go @@ -1340,13 +1340,13 @@ func TestOpenRestyProxyRequestBufferingDefaultsToOff(t *testing.T) { func setupServiceTestDB(t *testing.T) { t.Helper() - nodeAgentTokenCache.reset() + nodeAccessTokenCache.reset() common.SQLitePath = filepath.Join(t.TempDir(), "service.db") if err := model.InitDB(); err != nil { t.Fatalf("failed to init db: %v", err) } t.Cleanup(func() { - nodeAgentTokenCache.reset() + nodeAccessTokenCache.reset() if err := model.CloseDB(); err != nil { t.Fatalf("failed to close db: %v", err) } diff --git a/openflare_server/service/node.go b/openflare_server/service/node.go index 523c283b..6c4b41bf 100644 --- a/openflare_server/service/node.go +++ b/openflare_server/service/node.go @@ -57,9 +57,9 @@ type NodeBootstrapView struct { } type AgentRegistrationResponse struct { - NodeID string `json:"node_id"` - AgentToken string `json:"agent_token"` - Name string `json:"name"` + NodeID string `json:"node_id"` + AccessToken string `json:"access_token"` + Name string `json:"name"` } func CreateNode(input NodeInput) (*NodeView, error) { @@ -76,8 +76,8 @@ func CreateNode(input NodeInput) (*NodeView, error) { GeoLatitude: geoLatitude, GeoLongitude: geoLongitude, GeoManualOverride: geoManualOverride, - AgentVersion: "", - NginxVersion: "", + Version: "", + ExtVersion: "", Status: NodeStatusPending, AutoUpdateEnabled: input.AutoUpdateEnabled, NodeType: normalizeNodeType(input.NodeType), @@ -86,7 +86,7 @@ func CreateNode(input NodeInput) (*NodeView, error) { if err != nil { return nil, err } - node.AgentToken, err = newRandomToken() + node.AccessToken, err = newRandomToken() if err != nil { return nil, err } @@ -110,7 +110,7 @@ func CreateNode(input NodeInput) (*NodeView, error) { } return nil, err } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) slog.Info("node created", "name", node.Name, "node_id", node.NodeID) return buildNodeView(node), nil } @@ -150,7 +150,7 @@ func UpdateNode(id uint, input NodeInput) (*NodeView, error) { if err = node.Update(); err != nil { return nil, err } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) slog.Info("node updated", "name", node.Name, "node_id", node.NodeID) return buildNodeView(node), nil } @@ -164,7 +164,7 @@ func DeleteNode(id uint) error { if err := node.Delete(); err != nil { return err } - invalidateAgentTokenCache(node.AgentToken) + invalidateAccessTokenCache(node.AccessToken) DisconnectAgentWSClient(node.NodeID) return nil } @@ -206,7 +206,7 @@ func RequestNodeAgentUpdate(id uint, input NodeAgentUpdateInput) (*NodeView, err if err = model.DB.Model(node).Select("update_requested", "update_channel", "update_tag").Updates(node).Error; err != nil { return nil, err } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) if SendAgentWSSettings(node.NodeID, buildAgentSettings(node, true, channel.String(), tagName, node.RestartOpenrestyRequested)) { slog.Debug("agent manual update pushed via ws", "node_id", node.NodeID, "channel", channel.String(), "tag", tagName) } else { @@ -225,7 +225,7 @@ func RequestNodeOpenrestyRestart(id uint) (*NodeView, error) { if err = model.DB.Model(node).Select("restart_openresty_requested").Updates(node).Error; err != nil { return nil, err } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) slog.Info("openresty restart requested", "node_id", node.NodeID, "name", node.Name) return buildNodeView(node), nil } @@ -246,12 +246,12 @@ func RequestNodeForceSync(id uint) (*NodeView, error) { return buildNodeView(node), nil } -func AuthenticateAgentToken(token string) (*model.Node, error) { +func AuthenticateAccessToken(token string) (*model.Node, error) { token = strings.TrimSpace(token) if token == "" { return nil, errors.New("缺少 Agent Token") } - return authenticateAgentTokenWithCache(token) + return authenticateAccessTokenWithCache(token) } func ValidateDiscoveryToken(token string) error { @@ -323,12 +323,12 @@ func buildNodeView(node *model.Node) *NodeView { GeoLatitude: node.GeoLatitude, GeoLongitude: node.GeoLongitude, GeoManualOverride: node.GeoManualOverride, - AgentToken: node.AgentToken, + AccessToken: node.AccessToken, UpdateChannel: strings.TrimSpace(node.UpdateChannel), UpdateTag: strings.TrimSpace(node.UpdateTag), RestartOpenrestyRequested: node.RestartOpenrestyRequested, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + Version: node.Version, + ExtVersion: node.ExtVersion, OpenrestyStatus: normalizeOpenrestyStatus(node.OpenrestyStatus), OpenrestyMessage: strings.TrimSpace(node.OpenrestyMessage), Status: status, @@ -353,10 +353,8 @@ func buildNodeView(node *model.Node) *NodeView { view.RelayClientAccessAddr = node.RelayClientAccessAddr view.RelayClientProxyURL = node.RelayClientProxyURL view.RelayStatus = node.RelayStatus - view.RelayFrpVersion = node.RelayFrpVersion - view.RelayVersion = node.RelayVersion - view.RelayFrpsConnections = node.RelayFrpsConnections - view.RelayFrpsProxyCount = node.RelayFrpsProxyCount + view.Version = node.Version + view.ExtVersion = node.ExtVersion return view } @@ -455,7 +453,7 @@ func isPublicNodeIP(raw string) bool { } func buildNodeAgentReleaseView(node *model.Node, release *githubReleaseResponse, channel ReleaseChannel) *NodeAgentReleaseInfo { - currentVersion := strings.TrimSpace(node.AgentVersion) + currentVersion := strings.TrimSpace(node.Version) view := &NodeAgentReleaseInfo{ CurrentVersion: currentVersion, Channel: channel.String(), @@ -475,7 +473,7 @@ func buildNodeAgentReleaseView(node *model.Node, release *githubReleaseResponse, return view } -func RegisterNodeWithAgentToken(node *model.Node, payload AgentNodePayload) (*AgentRegistrationResponse, error) { +func RegisterNodeWithAccessToken(node *model.Node, payload AgentNodePayload) (*AgentRegistrationResponse, error) { payload = normalizeAgentNodePayload(payload) if node == nil { return nil, errors.New("节点不存在") @@ -487,12 +485,12 @@ func RegisterNodeWithAgentToken(node *model.Node, payload AgentNodePayload) (*Ag if err := node.Update(); err != nil { return nil, err } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) slog.Info("agent register succeeded on reserved node", "node_id", node.NodeID, "name", node.Name) return &AgentRegistrationResponse{ - NodeID: node.NodeID, - AgentToken: node.AgentToken, - Name: node.Name, + NodeID: node.NodeID, + AccessToken: node.AccessToken, + Name: node.Name, }, nil } @@ -514,9 +512,9 @@ func RegisterNodeWithDiscovery(payload AgentNodePayload) (*AgentRegistrationResp nodeName = nodeID } node := &model.Node{ - NodeID: nodeID, - Name: nodeName, - AgentToken: agentToken, + NodeID: nodeID, + Name: nodeName, + AccessToken: agentToken, } applyNodeRuntime(node, payload, false) if err = node.Insert(); err != nil { @@ -525,20 +523,20 @@ func RegisterNodeWithDiscovery(payload AgentNodePayload) (*AgentRegistrationResp } return nil, err } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) slog.Info("agent discovery register succeeded", "node_id", node.NodeID, "name", node.Name) return &AgentRegistrationResponse{ - NodeID: node.NodeID, - AgentToken: node.AgentToken, - Name: node.Name, + NodeID: node.NodeID, + AccessToken: node.AccessToken, + Name: node.Name, }, nil } func normalizeAgentNodePayload(payload AgentNodePayload) AgentNodePayload { payload.Name = strings.TrimSpace(payload.Name) payload.IP = strings.TrimSpace(payload.IP) - payload.AgentVersion = strings.TrimSpace(payload.AgentVersion) - payload.NginxVersion = strings.TrimSpace(payload.NginxVersion) + payload.Version = strings.TrimSpace(payload.Version) + payload.ExtVersion = strings.TrimSpace(payload.ExtVersion) payload.CurrentVersion = strings.TrimSpace(payload.CurrentVersion) payload.LastError = truncateForDatabase(payload.LastError, 16000) payload.OpenrestyStatus = normalizeOpenrestyStatus(payload.OpenrestyStatus) @@ -553,8 +551,8 @@ func validateAgentNodePayload(payload AgentNodePayload) error { if net.ParseIP(payload.IP) == nil { return errors.New("ip 格式无效") } - if payload.AgentVersion == "" { - return errors.New("agent_version 不能为空") + if payload.Version == "" { + return errors.New("version 不能为空") } return nil } @@ -568,8 +566,8 @@ func applyNodeRuntime(node *model.Node, payload AgentNodePayload, preserveName b if !node.IPManualOverride { node.IP = strings.TrimSpace(payload.IP) } - node.AgentVersion = strings.TrimSpace(payload.AgentVersion) - node.NginxVersion = strings.TrimSpace(payload.NginxVersion) + node.Version = strings.TrimSpace(payload.Version) + node.ExtVersion = strings.TrimSpace(payload.ExtVersion) node.OpenrestyStatus = normalizeOpenrestyStatus(payload.OpenrestyStatus) node.OpenrestyMessage = truncateForDatabase(payload.OpenrestyMessage, 16000) node.Status = NodeStatusOnline diff --git a/openflare_server/service/node_agent_token_cache.go b/openflare_server/service/node_agent_token_cache.go index 08a12496..04edf307 100644 --- a/openflare_server/service/node_agent_token_cache.go +++ b/openflare_server/service/node_agent_token_cache.go @@ -20,31 +20,31 @@ type cachedAgentNode struct { expiresAt time.Time } -type cachedMissingAgentToken struct { +type cachedMissingAccessToken struct { expiresAt time.Time } type agentTokenAuthCache struct { positive *ristretto.Cache[string, cachedAgentNode] - negative *ristretto.Cache[string, cachedMissingAgentToken] + negative *ristretto.Cache[string, cachedMissingAccessToken] now func() time.Time loadNodeByToken func(string) (*model.Node, error) } -var nodeAgentTokenCache = newAgentTokenAuthCache() +var nodeAccessTokenCache = newAccessTokenAuthCache() -func newAgentTokenAuthCache() *agentTokenAuthCache { +func newAccessTokenAuthCache() *agentTokenAuthCache { return &agentTokenAuthCache{ - positive: mustNewAgentTokenPositiveCache(), - negative: mustNewAgentTokenNegativeCache(), + positive: mustNewAccessTokenPositiveCache(), + negative: mustNewAccessTokenNegativeCache(), now: time.Now, loadNodeByToken: func(token string) (*model.Node, error) { - return model.GetNodeByAgentToken(token) + return model.GetNodeByAccessToken(token) }, } } -func mustNewAgentTokenPositiveCache() *ristretto.Cache[string, cachedAgentNode] { +func mustNewAccessTokenPositiveCache() *ristretto.Cache[string, cachedAgentNode] { cache, err := ristretto.NewCache(&ristretto.Config[string, cachedAgentNode]{ NumCounters: 1e5, MaxCost: 2e4, @@ -56,8 +56,8 @@ func mustNewAgentTokenPositiveCache() *ristretto.Cache[string, cachedAgentNode] return cache } -func mustNewAgentTokenNegativeCache() *ristretto.Cache[string, cachedMissingAgentToken] { - cache, err := ristretto.NewCache(&ristretto.Config[string, cachedMissingAgentToken]{ +func mustNewAccessTokenNegativeCache() *ristretto.Cache[string, cachedMissingAccessToken] { + cache, err := ristretto.NewCache(&ristretto.Config[string, cachedMissingAccessToken]{ NumCounters: 1e5, MaxCost: agentTokenNegativeCacheCap, BufferItems: 64, @@ -130,7 +130,7 @@ func (c *agentTokenAuthCache) storeMissing(token string, expiresAt time.Time) { return } c.positive.Del(token) - c.negative.Set(token, cachedMissingAgentToken{ + c.negative.Set(token, cachedMissingAccessToken{ expiresAt: expiresAt, }, 1) c.negative.Wait() @@ -157,21 +157,21 @@ func cloneCachedNode(node *model.Node) *model.Node { return &cloned } -func authenticateAgentTokenWithCache(token string) (*model.Node, error) { - return nodeAgentTokenCache.authenticate(token) +func authenticateAccessTokenWithCache(token string) (*model.Node, error) { + return nodeAccessTokenCache.authenticate(token) } -func refreshAgentTokenCache(node *model.Node) { +func refreshAccessTokenCache(node *model.Node) { if node == nil { return } - nodeAgentTokenCache.storeNode( - node.AgentToken, + nodeAccessTokenCache.storeNode( + node.AccessToken, node, - nodeAgentTokenCache.now().Add(agentTokenPositiveCacheTTL), + nodeAccessTokenCache.now().Add(agentTokenPositiveCacheTTL), ) } -func invalidateAgentTokenCache(token string) { - nodeAgentTokenCache.invalidate(token) +func invalidateAccessTokenCache(token string) { + nodeAccessTokenCache.invalidate(token) } diff --git a/openflare_server/service/node_agent_token_cache_test.go b/openflare_server/service/node_agent_token_cache_test.go index cf56ed58..2b4aeb1a 100644 --- a/openflare_server/service/node_agent_token_cache_test.go +++ b/openflare_server/service/node_agent_token_cache_test.go @@ -10,8 +10,8 @@ import ( "gorm.io/gorm" ) -func TestAgentTokenAuthCacheUsesPositiveCacheUntilLogicalExpiry(t *testing.T) { - cache := newAgentTokenAuthCache() +func TestAccessTokenAuthCacheUsesPositiveCacheUntilLogicalExpiry(t *testing.T) { + cache := newAccessTokenAuthCache() cache.reset() baseTime := time.Date(2026, 3, 14, 16, 0, 0, 0, time.UTC) currentTime := baseTime @@ -23,9 +23,9 @@ func TestAgentTokenAuthCacheUsesPositiveCacheUntilLogicalExpiry(t *testing.T) { cache.loadNodeByToken = func(token string) (*model.Node, error) { loadCount++ return &model.Node{ - NodeID: fmt.Sprintf("node-%d", loadCount), - Name: "edge", - AgentToken: token, + NodeID: fmt.Sprintf("node-%d", loadCount), + Name: "edge", + AccessToken: token, }, nil } @@ -61,8 +61,8 @@ func TestAgentTokenAuthCacheUsesPositiveCacheUntilLogicalExpiry(t *testing.T) { } } -func TestAgentTokenAuthCacheRefreshesAfterMissingEntryExpires(t *testing.T) { - cache := newAgentTokenAuthCache() +func TestAccessTokenAuthCacheRefreshesAfterMissingEntryExpires(t *testing.T) { + cache := newAccessTokenAuthCache() cache.reset() baseTime := time.Date(2026, 3, 14, 16, 30, 0, 0, time.UTC) currentTime := baseTime @@ -77,9 +77,9 @@ func TestAgentTokenAuthCacheRefreshesAfterMissingEntryExpires(t *testing.T) { return nil, gorm.ErrRecordNotFound } return &model.Node{ - NodeID: "node-recovered", - Name: "edge", - AgentToken: token, + NodeID: "node-recovered", + Name: "edge", + AccessToken: token, }, nil } diff --git a/openflare_server/service/node_observability.go b/openflare_server/service/node_observability.go index 11306591..a97e32dc 100644 --- a/openflare_server/service/node_observability.go +++ b/openflare_server/service/node_observability.go @@ -101,6 +101,7 @@ func GetNodeObservability(id uint, query NodeObservabilityQuery) (*NodeObservabi if err != nil { return nil, err } + trendOpenresty, _ := model.ListNodeObservationOpenresty(node.NodeID, now.Add(-24*time.Hour), 0) trendReports, err := model.ListNodeRequestReports(node.NodeID, now.Add(-24*time.Hour), 0) if err != nil { return nil, err @@ -124,21 +125,31 @@ func GetNodeObservability(id uint, query NodeObservabilityQuery) (*NodeObservabi Trends: NodeObservabilityTrends{ Traffic24h: buildTrafficTrendPoints(now, trendReports), Capacity24h: buildCapacityTrendPoints(now, trendSnapshots), - Network24h: buildNetworkTrendPoints(now, trendSnapshots), + Network24h: buildNetworkTrendPoints(now, trendSnapshots, trendOpenresty), DiskIO24h: buildDiskIOTrendPoints(now, trendSnapshots), }, } if node.NodeType == "tunnel_relay" { - view.RelayDashboard = buildRelayDashboardSnapshot(node) + frpsObs, _ := model.ListNodeObservationFrps(node.NodeID, time.Time{}, 1) + var latestFrps *model.NodeObservationFrps + if len(frpsObs) > 0 { + latestFrps = frpsObs[0] + } + view.RelayDashboard = buildRelayDashboardSnapshot(node, latestFrps) } return view, nil } -func buildRelayDashboardSnapshot(node *model.Node) *RelayDashboardSnapshot { +func buildRelayDashboardSnapshot(node *model.Node, obs *model.NodeObservationFrps) *RelayDashboardSnapshot { if node == nil { return nil } - totalProxies := node.RelayFrpsProxyCount + totalProxies := 0 + totalConnections := 0 + if obs != nil { + totalProxies = obs.FrpsProxyCount + totalConnections = obs.FrpsConnections + } if totalProxies < 0 { totalProxies = 0 } @@ -151,7 +162,7 @@ func buildRelayDashboardSnapshot(node *model.Node) *RelayDashboardSnapshot { OnlineProxies: onlineProxies, OfflineProxies: totalProxies - onlineProxies, Proxies: []RelayProxyStat{}, - TotalConnections: maxInt(node.RelayFrpsConnections, 0), + TotalConnections: maxInt(totalConnections, 0), ClientCounts: 0, } } diff --git a/openflare_server/service/node_update_test.go b/openflare_server/service/node_update_test.go index e5f6b842..cbf07de8 100644 --- a/openflare_server/service/node_update_test.go +++ b/openflare_server/service/node_update_test.go @@ -167,9 +167,9 @@ func TestHeartbeatNodeReturnsPreviewUpdateSettings(t *testing.T) { NodeID: "node-preview-1", Name: "preview-edge-1", IP: "10.0.0.8", - AgentToken: "agent-token", - AgentVersion: "v0.4.0", - NginxVersion: "1.27.1.2", + AccessToken: "agent-token", + Version: "v0.4.0", + ExtVersion: "1.27.1.2", Status: NodeStatusOnline, UpdateRequested: true, UpdateChannel: "preview", @@ -196,8 +196,8 @@ func TestHeartbeatNodeReturnsPreviewUpdateSettings(t *testing.T) { NodeID: node.NodeID, Name: node.Name, IP: node.IP, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + Version: node.Version, + ExtVersion: node.ExtVersion, OpenrestyStatus: OpenrestyStatusUnhealthy, OpenrestyMessage: "port 80 already allocated", }) @@ -381,24 +381,24 @@ func TestHeartbeatNodeResolvesGeoMetadataFromIPWhenNotManuallyOverridden(t *test }) node := &model.Node{ - NodeID: "node-geo-auto", - Name: "geo-auto", - IP: "10.0.0.8", - AgentToken: "agent-token", - AgentVersion: "v0.4.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-geo-auto", + Name: "geo-auto", + IP: "10.0.0.8", + AccessToken: "agent-token", + Version: "v0.4.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) } resp, err := HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: "8.8.8.8", - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: "8.8.8.8", + Version: node.Version, + ExtVersion: node.ExtVersion, }) if err != nil { t.Fatalf("expected heartbeat to succeed: %v", err) @@ -430,9 +430,9 @@ func TestHeartbeatNodePreservesManualGeoOverride(t *testing.T) { GeoLatitude: &latitude, GeoLongitude: &longitude, GeoManualOverride: true, - AgentToken: "agent-token", - AgentVersion: "v0.4.0", - NginxVersion: "1.27.1.2", + AccessToken: "agent-token", + Version: "v0.4.0", + ExtVersion: "1.27.1.2", Status: NodeStatusOnline, } if err := node.Insert(); err != nil { @@ -440,11 +440,11 @@ func TestHeartbeatNodePreservesManualGeoOverride(t *testing.T) { } resp, err := HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: "8.8.8.8", - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: "8.8.8.8", + Version: node.Version, + ExtVersion: node.ExtVersion, }) if err != nil { t.Fatalf("expected heartbeat to succeed: %v", err) @@ -465,9 +465,9 @@ func TestHeartbeatNodePreservesManualIPOverride(t *testing.T) { Name: "ip-manual", IP: "203.0.113.10", IPManualOverride: true, - AgentToken: "agent-token", - AgentVersion: "v0.4.0", - NginxVersion: "1.27.1.2", + AccessToken: "agent-token", + Version: "v0.4.0", + ExtVersion: "1.27.1.2", Status: NodeStatusOnline, } if err := node.Insert(); err != nil { @@ -478,8 +478,8 @@ func TestHeartbeatNodePreservesManualIPOverride(t *testing.T) { NodeID: node.NodeID, Name: node.Name, IP: "10.0.0.8", - AgentVersion: "v0.5.0", - NginxVersion: "1.27.1.3", + Version: "v0.5.0", + ExtVersion: "1.27.1.3", OpenrestyStatus: OpenrestyStatusHealthy, }) if err != nil { @@ -488,7 +488,7 @@ func TestHeartbeatNodePreservesManualIPOverride(t *testing.T) { if resp.Node.IP != "203.0.113.10" { t.Fatalf("expected manual ip to be preserved, got %s", resp.Node.IP) } - if resp.Node.AgentVersion != "v0.5.0" || resp.Node.NginxVersion != "1.27.1.3" { + if resp.Node.Version != "v0.5.0" || resp.Node.ExtVersion != "1.27.1.3" { t.Fatalf("expected runtime metadata to update despite locked ip, got %+v", resp.Node) } @@ -505,24 +505,24 @@ func TestHeartbeatNodeUpdatesIPWhenManualOverrideDisabled(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-ip-auto", - Name: "ip-auto", - IP: "10.0.0.8", - AgentToken: "agent-token", - AgentVersion: "v0.4.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-ip-auto", + Name: "ip-auto", + IP: "10.0.0.8", + AccessToken: "agent-token", + Version: "v0.4.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) } resp, err := HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: "8.8.8.8", - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: "8.8.8.8", + Version: node.Version, + ExtVersion: node.ExtVersion, }) if err != nil { t.Fatalf("expected heartbeat to succeed: %v", err) @@ -557,11 +557,11 @@ func TestUpdateNodeCanLockAndUnlockManualIP(t *testing.T) { t.Fatalf("failed to reload node: %v", err) } if _, err = HeartbeatNode(stored, AgentNodePayload{ - NodeID: stored.NodeID, - Name: stored.Name, - IP: "8.8.8.8", - AgentVersion: "v0.5.0", - NginxVersion: "1.27.1.3", + NodeID: stored.NodeID, + Name: stored.Name, + IP: "8.8.8.8", + Version: "v0.5.0", + ExtVersion: "1.27.1.3", }); err != nil { t.Fatalf("expected heartbeat to succeed: %v", err) } @@ -591,11 +591,11 @@ func TestUpdateNodeCanLockAndUnlockManualIP(t *testing.T) { t.Fatalf("failed to reload unlocked node: %v", err) } if _, err = HeartbeatNode(unlockedStored, AgentNodePayload{ - NodeID: unlockedStored.NodeID, - Name: unlockedStored.Name, - IP: "8.8.4.4", - AgentVersion: "v0.5.1", - NginxVersion: "1.27.1.4", + NodeID: unlockedStored.NodeID, + Name: unlockedStored.Name, + IP: "8.8.4.4", + Version: "v0.5.1", + ExtVersion: "1.27.1.4", }); err != nil { t.Fatalf("expected heartbeat to succeed after unlock: %v", err) } @@ -666,25 +666,25 @@ func TestListNodeViewsIncludesLatestApplyLogsForMultipleNodes(t *testing.T) { now := time.Now() nodes := []*model.Node{ { - NodeID: "node-a", - Name: "edge-a", - IP: "10.0.0.11", - GeoName: "Shanghai", - AgentToken: "token-a", - AgentVersion: "v0.5.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, - LastSeenAt: now, + NodeID: "node-a", + Name: "edge-a", + IP: "10.0.0.11", + GeoName: "Shanghai", + AccessToken: "token-a", + Version: "v0.5.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, + LastSeenAt: now, }, { - NodeID: "node-b", - Name: "edge-b", - IP: "10.0.0.12", - AgentToken: "token-b", - AgentVersion: "v0.5.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, - LastSeenAt: now, + NodeID: "node-b", + Name: "edge-b", + IP: "10.0.0.12", + AccessToken: "token-b", + Version: "v0.5.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, + LastSeenAt: now, }, } for _, node := range nodes { @@ -732,8 +732,8 @@ func TestCollectNodeHeartbeatChangesOnlyReturnsChangedFields(t *testing.T) { before := &model.Node{ Name: "edge-1", IP: "10.0.0.8", - AgentVersion: "v0.5.0", - NginxVersion: "1.27.1.2", + Version: "v0.5.0", + ExtVersion: "1.27.1.2", OpenrestyStatus: OpenrestyStatusHealthy, OpenrestyMessage: "", Status: NodeStatusOnline, @@ -748,8 +748,8 @@ func TestCollectNodeHeartbeatChangesOnlyReturnsChangedFields(t *testing.T) { after := &model.Node{ Name: "edge-1", IP: "10.0.0.8", - AgentVersion: "v0.5.0", - NginxVersion: "1.27.1.2", + Version: "v0.5.0", + ExtVersion: "1.27.1.2", OpenrestyStatus: OpenrestyStatusHealthy, OpenrestyMessage: "", Status: NodeStatusOnline, @@ -790,14 +790,14 @@ func TestListNodeViewsDoesNotPersistComputedStatus(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-offline-view", - Name: "edge-offline", - IP: "10.0.0.21", - AgentToken: "token-offline", - AgentVersion: "v0.5.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, - LastSeenAt: time.Now().Add(-common.NodeOfflineThreshold - time.Minute), + NodeID: "node-offline-view", + Name: "edge-offline", + IP: "10.0.0.21", + AccessToken: "token-offline", + Version: "v0.5.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, + LastSeenAt: time.Now().Add(-common.NodeOfflineThreshold - time.Minute), } if err := node.Insert(); err != nil { t.Fatalf("failed to insert node: %v", err) @@ -868,24 +868,24 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) { }) node := &model.Node{ - NodeID: "node-observe-1", - Name: "observe-edge-1", - IP: "10.0.0.31", - AgentToken: "token-observe", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-observe-1", + Name: "observe-edge-1", + IP: "10.0.0.31", + AccessToken: "token-observe", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) } _, err := HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: node.IP, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + Version: node.Version, + ExtVersion: node.ExtVersion, Profile: &AgentNodeSystemProfile{ Hostname: "observe-edge-1", OSName: "Ubuntu", @@ -900,17 +900,16 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) { ReportedAtUnix: time.Now().Add(-time.Minute).Unix(), }, Snapshot: &AgentNodeMetricSnapshot{ - CapturedAtUnix: time.Now().Add(-30 * time.Second).Unix(), - CPUUsagePercent: 42.5, - MemoryUsedBytes: 8 * 1024 * 1024 * 1024, - MemoryTotalBytes: 16 * 1024 * 1024 * 1024, - StorageUsedBytes: 70 * 1024 * 1024 * 1024, - StorageTotalBytes: 200 * 1024 * 1024 * 1024, - DiskReadBytes: 1024, - DiskWriteBytes: 2048, - NetworkRxBytes: 4096, - NetworkTxBytes: 8192, - OpenrestyConnections: 128, + CapturedAtUnix: time.Now().Add(-30 * time.Second).Unix(), + CPUUsagePercent: 42.5, + MemoryUsedBytes: 8 * 1024 * 1024 * 1024, + MemoryTotalBytes: 16 * 1024 * 1024 * 1024, + StorageUsedBytes: 70 * 1024 * 1024 * 1024, + StorageTotalBytes: 200 * 1024 * 1024 * 1024, + DiskReadBytes: 1024, + DiskWriteBytes: 2048, + NetworkRxBytes: 4096, + NetworkTxBytes: 8192, }, TrafficReport: &AgentNodeTrafficReport{ WindowStartedAtUnix: time.Now().Add(-time.Minute).Unix(), @@ -966,7 +965,7 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) { if err != nil { t.Fatalf("expected node snapshots query to succeed: %v", err) } - if len(snapshots) != 1 || snapshots[0].OpenrestyConnections != 128 { + if len(snapshots) != 1 { t.Fatalf("unexpected metric snapshots: %+v", snapshots) } @@ -1021,13 +1020,13 @@ func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) { }) node := &model.Node{ - NodeID: "node-observe-buffered", - Name: "observe-buffered-edge", - IP: "10.0.0.32", - AgentToken: "token-observe-buffered", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-observe-buffered", + Name: "observe-buffered-edge", + IP: "10.0.0.32", + AccessToken: "token-observe-buffered", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) @@ -1035,11 +1034,11 @@ func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) { now := time.Now().UTC() _, err := HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: node.IP, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + Version: node.Version, + ExtVersion: node.ExtVersion, Snapshot: &AgentNodeMetricSnapshot{ CapturedAtUnix: now.Unix(), CPUUsagePercent: 25, @@ -1124,11 +1123,11 @@ func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) { } _, err = HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: node.IP, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + Version: node.Version, + ExtVersion: node.ExtVersion, BufferedObservability: []AgentBufferedObservabilityRecord{ { WindowStartedAtUnix: now.Add(-2 * time.Minute).Unix(), @@ -1196,13 +1195,13 @@ func TestListAccessLogsUsesPagination(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-access-log-page", - Name: "access-log-edge", - IP: "10.0.0.40", - AgentToken: "token-access-log-page", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-access-log-page", + Name: "access-log-edge", + IP: "10.0.0.40", + AccessToken: "token-access-log-page", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) @@ -1281,24 +1280,24 @@ func TestHeartbeatNodeResolvesMissingHealthEvents(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-event-1", - Name: "event-edge-1", - IP: "10.0.0.41", - AgentToken: "token-event", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-event-1", + Name: "event-edge-1", + IP: "10.0.0.41", + AccessToken: "token-event", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to seed node: %v", err) } _, err := HeartbeatNode(node, AgentNodePayload{ - NodeID: node.NodeID, - Name: node.Name, - IP: node.IP, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + Version: node.Version, + ExtVersion: node.ExtVersion, HealthEvents: []AgentNodeHealthEvent{ { EventType: "sync_error", @@ -1316,8 +1315,8 @@ func TestHeartbeatNodeResolvesMissingHealthEvents(t *testing.T) { NodeID: node.NodeID, Name: node.Name, IP: node.IP, - AgentVersion: node.AgentVersion, - NginxVersion: node.NginxVersion, + Version: node.Version, + ExtVersion: node.ExtVersion, HealthEvents: []AgentNodeHealthEvent{}, }) if err != nil { @@ -1345,13 +1344,13 @@ func TestGetNodeObservability(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-observability-query", - Name: "query-edge", - IP: "10.0.0.61", - AgentToken: "token-query", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-observability-query", + Name: "query-edge", + IP: "10.0.0.61", + AccessToken: "token-query", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to insert node: %v", err) @@ -1377,8 +1376,6 @@ func TestGetNodeObservability(t *testing.T) { DiskWriteBytes: 0, NetworkRxBytes: 2048, NetworkTxBytes: 4096, - OpenrestyRxBytes: 8192, - OpenrestyTxBytes: 16384, }).Insert(); err != nil { t.Fatalf("failed to insert node metric baseline snapshot: %v", err) } @@ -1394,8 +1391,6 @@ func TestGetNodeObservability(t *testing.T) { DiskWriteBytes: 2048, NetworkRxBytes: 4096, NetworkTxBytes: 8192, - OpenrestyRxBytes: 16384, - OpenrestyTxBytes: 32768, }).Insert(); err != nil { t.Fatalf("failed to insert node metric snapshot: %v", err) } @@ -1453,7 +1448,7 @@ func TestGetNodeObservability(t *testing.T) { if view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1].RequestCount != 123 { t.Fatalf("unexpected traffic trend tail: %+v", view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1]) } - if view.Trends.Network24h[len(view.Trends.Network24h)-1].OpenrestyTxBytes != 32768 { + if view.Trends.Network24h[len(view.Trends.Network24h)-1].NetworkTxBytes != 8192 { t.Fatalf("unexpected network trend tail: %+v", view.Trends.Network24h[len(view.Trends.Network24h)-1]) } if view.Trends.DiskIO24h[len(view.Trends.DiskIO24h)-1].DiskWriteBytes != 2048 { @@ -1477,13 +1472,13 @@ func TestGetNodeObservabilityAllowsMissingProfile(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-observability-empty", - Name: "empty-edge", - IP: "10.0.0.62", - AgentToken: "token-empty", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-observability-empty", + Name: "empty-edge", + IP: "10.0.0.62", + AccessToken: "token-empty", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to insert node: %v", err) @@ -1508,13 +1503,13 @@ func TestCleanupNodeHealthEvents(t *testing.T) { setupServiceTestDB(t) node := &model.Node{ - NodeID: "node-health-cleanup", - Name: "health-cleanup-edge", - IP: "10.0.0.72", - AgentToken: "token-health-cleanup", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", - Status: NodeStatusOnline, + NodeID: "node-health-cleanup", + Name: "health-cleanup-edge", + IP: "10.0.0.72", + AccessToken: "token-health-cleanup", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", + Status: NodeStatusOnline, } if err := node.Insert(); err != nil { t.Fatalf("failed to insert node: %v", err) @@ -1586,9 +1581,9 @@ func TestGetDashboardOverview(t *testing.T) { Name: "edge-a", IP: "10.0.0.71", GeoName: "Shanghai", - AgentToken: "token-a", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", + AccessToken: "token-a", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", OpenrestyStatus: OpenrestyStatusHealthy, Status: NodeStatusOnline, CurrentVersion: "20260314-001", @@ -1599,9 +1594,9 @@ func TestGetDashboardOverview(t *testing.T) { Name: "edge-b", IP: "10.0.0.72", GeoName: "San Francisco", - AgentToken: "token-b", - AgentVersion: "v0.6.0", - NginxVersion: "1.27.1.2", + AccessToken: "token-b", + Version: "v0.6.0", + ExtVersion: "1.27.1.2", OpenrestyStatus: OpenrestyStatusUnhealthy, Status: NodeStatusOnline, CurrentVersion: "20260313-001", @@ -1626,8 +1621,6 @@ func TestGetDashboardOverview(t *testing.T) { DiskWriteBytes: 0, NetworkRxBytes: 100, NetworkTxBytes: 150, - OpenrestyRxBytes: 300, - OpenrestyTxBytes: 450, }).Insert(); err != nil { t.Fatalf("failed to insert node a baseline metric snapshot: %v", err) } @@ -1643,8 +1636,6 @@ func TestGetDashboardOverview(t *testing.T) { DiskWriteBytes: 150, NetworkRxBytes: 300, NetworkTxBytes: 500, - OpenrestyRxBytes: 700, - OpenrestyTxBytes: 900, }).Insert(); err != nil { t.Fatalf("failed to insert node a metric snapshot: %v", err) } @@ -1659,9 +1650,7 @@ func TestGetDashboardOverview(t *testing.T) { DiskReadBytes: 0, DiskWriteBytes: 0, NetworkRxBytes: 200, - NetworkTxBytes: 300, - OpenrestyRxBytes: 500, - OpenrestyTxBytes: 700, + NetworkTxBytes: 400, }).Insert(); err != nil { t.Fatalf("failed to insert node b baseline metric snapshot: %v", err) } @@ -1676,9 +1665,7 @@ func TestGetDashboardOverview(t *testing.T) { DiskReadBytes: 200, DiskWriteBytes: 400, NetworkRxBytes: 600, - NetworkTxBytes: 900, - OpenrestyRxBytes: 1200, - OpenrestyTxBytes: 1600, + NetworkTxBytes: 800, }).Insert(); err != nil { t.Fatalf("failed to insert node b metric snapshot: %v", err) } @@ -1785,7 +1772,7 @@ func TestGetDashboardOverview(t *testing.T) { if view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1].RequestCount != 900 { t.Fatalf("unexpected dashboard traffic trend tail: %+v", view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1]) } - if view.Trends.Network24h[len(view.Trends.Network24h)-1].OpenrestyRxBytes != 1900 { + if view.Trends.Network24h[len(view.Trends.Network24h)-1].NetworkRxBytes != 900 { t.Fatalf("unexpected dashboard network trend tail: %+v", view.Trends.Network24h[len(view.Trends.Network24h)-1]) } if view.Trends.DiskIO24h[len(view.Trends.DiskIO24h)-1].DiskWriteBytes != 550 { diff --git a/openflare_server/service/observability.go b/openflare_server/service/observability.go index ccee23d8..0f6dc2d7 100644 --- a/openflare_server/service/observability.go +++ b/openflare_server/service/observability.go @@ -36,19 +36,23 @@ type AgentNodeSystemProfile struct { } type AgentNodeMetricSnapshot struct { - CapturedAtUnix int64 `json:"captured_at_unix"` - CPUUsagePercent float64 `json:"cpu_usage_percent"` - MemoryUsedBytes int64 `json:"memory_used_bytes"` - MemoryTotalBytes int64 `json:"memory_total_bytes"` - StorageUsedBytes int64 `json:"storage_used_bytes"` - StorageTotalBytes int64 `json:"storage_total_bytes"` - DiskReadBytes int64 `json:"disk_read_bytes"` - DiskWriteBytes int64 `json:"disk_write_bytes"` - NetworkRxBytes int64 `json:"network_rx_bytes"` - NetworkTxBytes int64 `json:"network_tx_bytes"` - OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` - OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` - OpenrestyConnections int64 `json:"openresty_connections"` + CapturedAtUnix int64 `json:"captured_at_unix"` + CPUUsagePercent float64 `json:"cpu_usage_percent"` + MemoryUsedBytes int64 `json:"memory_used_bytes"` + MemoryTotalBytes int64 `json:"memory_total_bytes"` + StorageUsedBytes int64 `json:"storage_used_bytes"` + StorageTotalBytes int64 `json:"storage_total_bytes"` + DiskReadBytes int64 `json:"disk_read_bytes"` + DiskWriteBytes int64 `json:"disk_write_bytes"` + NetworkRxBytes int64 `json:"network_rx_bytes"` + NetworkTxBytes int64 `json:"network_tx_bytes"` +} + +type AgentNodeOpenrestyObservation struct { + CapturedAtUnix int64 `json:"captured_at_unix"` + OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` + OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` + OpenrestyConnections int64 `json:"openresty_connections"` } type AgentNodeTrafficReport struct { @@ -71,10 +75,11 @@ type AgentNodeAccessLog struct { } type AgentBufferedObservabilityRecord struct { - WindowStartedAtUnix int64 `json:"window_started_at_unix"` - Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"` - TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"` - AccessLogs []AgentNodeAccessLog `json:"access_logs,omitempty"` + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"` + OpenrestyObservation *AgentNodeOpenrestyObservation `json:"openresty_observation,omitempty"` + TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"` + AccessLogs []AgentNodeAccessLog `json:"access_logs,omitempty"` } type AgentNodeHealthEvent struct { @@ -103,6 +108,9 @@ func persistHeartbeatObservability(nodeID string, payload AgentNodePayload, repo if err := persistNodeMetricSnapshot(tx, nodeID, payload.Snapshot, reportedAt); err != nil { return err } + if err := persistNodeOpenrestyObservation(tx, nodeID, payload.OpenrestyObservation, reportedAt); err != nil { + return err + } if err := persistNodeTrafficReport(tx, nodeID, payload.TrafficReport, reportedAt); err != nil { return err } @@ -125,6 +133,9 @@ func persistBufferedObservability(tx *gorm.DB, nodeID string, records []AgentBuf if err := persistNodeMetricSnapshot(tx, nodeID, record.Snapshot, reportedAt); err != nil { return err } + if err := persistNodeOpenrestyObservation(tx, nodeID, record.OpenrestyObservation, reportedAt); err != nil { + return err + } if err := persistNodeTrafficReport(tx, nodeID, record.TrafficReport, reportedAt); err != nil { return err } @@ -161,20 +172,17 @@ func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *AgentNodeMe return nil } record := &model.NodeMetricSnapshot{ - NodeID: nodeID, - CapturedAt: timeFromUnix(snapshot.CapturedAtUnix, reportedAt), - CPUUsagePercent: snapshot.CPUUsagePercent, - MemoryUsedBytes: snapshot.MemoryUsedBytes, - MemoryTotalBytes: snapshot.MemoryTotalBytes, - StorageUsedBytes: snapshot.StorageUsedBytes, - StorageTotalBytes: snapshot.StorageTotalBytes, - DiskReadBytes: snapshot.DiskReadBytes, - DiskWriteBytes: snapshot.DiskWriteBytes, - NetworkRxBytes: snapshot.NetworkRxBytes, - NetworkTxBytes: snapshot.NetworkTxBytes, - OpenrestyRxBytes: snapshot.OpenrestyRxBytes, - OpenrestyTxBytes: snapshot.OpenrestyTxBytes, - OpenrestyConnections: snapshot.OpenrestyConnections, + NodeID: nodeID, + CapturedAt: timeFromUnix(snapshot.CapturedAtUnix, reportedAt), + CPUUsagePercent: snapshot.CPUUsagePercent, + MemoryUsedBytes: snapshot.MemoryUsedBytes, + MemoryTotalBytes: snapshot.MemoryTotalBytes, + StorageUsedBytes: snapshot.StorageUsedBytes, + StorageTotalBytes: snapshot.StorageTotalBytes, + DiskReadBytes: snapshot.DiskReadBytes, + DiskWriteBytes: snapshot.DiskWriteBytes, + NetworkRxBytes: snapshot.NetworkRxBytes, + NetworkTxBytes: snapshot.NetworkTxBytes, } exists, err := model.NodeMetricSnapshotExists(tx, nodeID, record.CapturedAt) if err != nil { @@ -186,6 +194,20 @@ func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *AgentNodeMe return tx.Create(record).Error } +func persistNodeOpenrestyObservation(tx *gorm.DB, nodeID string, obs *AgentNodeOpenrestyObservation, reportedAt time.Time) error { + if obs == nil { + return nil + } + record := &model.NodeObservationOpenresty{ + NodeID: nodeID, + CapturedAt: timeFromUnix(obs.CapturedAtUnix, reportedAt), + OpenrestyRxBytes: obs.OpenrestyRxBytes, + OpenrestyTxBytes: obs.OpenrestyTxBytes, + OpenrestyConnections: obs.OpenrestyConnections, + } + return tx.Create(record).Error +} + func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *AgentNodeTrafficReport, reportedAt time.Time) error { if report == nil { return nil diff --git a/openflare_server/service/observability_trends.go b/openflare_server/service/observability_trends.go index 56274498..c110e294 100644 --- a/openflare_server/service/observability_trends.go +++ b/openflare_server/service/observability_trends.go @@ -110,7 +110,7 @@ func buildCapacityTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapsh return points } -func buildNetworkTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapshot) []NetworkTrendPoint { +func buildNetworkTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapshot, openrestyObs []*model.NodeObservationOpenresty) []NetworkTrendPoint { start := trendWindowStart(now) points := make([]NetworkTrendPoint, observabilityTrendBuckets) accumulators := make([]snapshotTrendAccumulator, observabilityTrendBuckets) @@ -126,13 +126,23 @@ func buildNetworkTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapsho } points[index].NetworkRxBytes += snapshot.NetworkRxBytes points[index].NetworkTxBytes += snapshot.NetworkTxBytes - points[index].OpenrestyRxBytes += snapshot.OpenrestyRxBytes - points[index].OpenrestyTxBytes += snapshot.OpenrestyTxBytes if snapshot.NodeID != "" { accumulators[index].nodes[snapshot.NodeID] = struct{}{} } } + for _, obs := range openrestyObs { + index, ok := trendBucketIndex(obs.CapturedAt, start) + if !ok { + continue + } + points[index].OpenrestyRxBytes += obs.OpenrestyRxBytes + points[index].OpenrestyTxBytes += obs.OpenrestyTxBytes + if obs.NodeID != "" { + accumulators[index].nodes[obs.NodeID] = struct{}{} + } + } + for index := range points { points[index].ReportedNodes = len(accumulators[index].nodes) } diff --git a/openflare_server/service/relay.go b/openflare_server/service/relay.go index c414fa12..153a6a20 100644 --- a/openflare_server/service/relay.go +++ b/openflare_server/service/relay.go @@ -11,8 +11,8 @@ import ( // RelayHeartbeatPayload is the payload sent by OpenFlareRelay in each heartbeat. type RelayHeartbeatPayload struct { - RelayVersion string `json:"relay_version"` - FrpVersion string `json:"frp_version"` + Version string `json:"version"` + ExtVersion string `json:"frp_version"` RelayStatus string `json:"relay_status"` FrpsConnCount int `json:"frps_connections"` FrpsProxyCount int `json:"frps_proxy_count"` @@ -50,8 +50,8 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear } slog.Debug("relay heartbeat received", "node_id", node.NodeID) - payload.RelayVersion = strings.TrimSpace(payload.RelayVersion) - payload.FrpVersion = strings.TrimSpace(payload.FrpVersion) + payload.Version = strings.TrimSpace(payload.Version) + payload.ExtVersion = strings.TrimSpace(payload.ExtVersion) payload.RelayStatus = normalizeRelayStatus(payload.RelayStatus) payload.Name = strings.TrimSpace(payload.Name) payload.IP = strings.TrimSpace(payload.IP) @@ -63,11 +63,10 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear } } now := time.Now() - appendRelayChange("relay_version", node.RelayVersion, payload.RelayVersion) - appendRelayChange("relay_frp_version", node.RelayFrpVersion, payload.FrpVersion) + appendRelayChange("version", node.Version, payload.Version) + appendRelayChange("ext_version", node.ExtVersion, payload.ExtVersion) 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 @@ -87,11 +86,10 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear } changes["status"] = NodeStatusOnline - node.RelayVersion = payload.RelayVersion - node.RelayFrpVersion = payload.FrpVersion + node.Version = payload.Version + node.ExtVersion = payload.ExtVersion node.RelayStatus = payload.RelayStatus - node.RelayFrpsConnections = payload.FrpsConnCount - node.RelayFrpsProxyCount = payload.FrpsProxyCount + node.LastSeenAt = now node.Status = NodeStatusOnline @@ -100,7 +98,7 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear return nil, fmt.Errorf("update relay heartbeat: %w", err) } } - refreshAgentTokenCache(node) + refreshAccessTokenCache(node) persistRelayHeartbeatObservability(node.NodeID, payload, node.LastSeenAt) return &RelayHeartbeatResponse{ @@ -115,6 +113,14 @@ func persistRelayHeartbeatObservability(nodeID string, payload RelayHeartbeatPay Snapshot: payload.Snapshot, HealthEvents: payload.HealthEvents, }, reportedAt) + + frpsObs := &model.NodeObservationFrps{ + NodeID: nodeID, + CapturedAt: reportedAt, + FrpsConnections: payload.FrpsConnCount, + FrpsProxyCount: payload.FrpsProxyCount, + } + _ = frpsObs.Insert() } func buildRelayConfig(node *model.Node) *RelayConfig { diff --git a/openflare_server/service/relay_test.go b/openflare_server/service/relay_test.go index 60cbb78d..af2144be 100644 --- a/openflare_server/service/relay_test.go +++ b/openflare_server/service/relay_test.go @@ -10,14 +10,14 @@ 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: "", + NodeID: "node-relay-observe", + Name: "relay-1", + IP: "", + AccessToken: "relay-token", + Status: NodeStatusPending, + NodeType: "tunnel_relay", + RelayStatus: "unknown", + Version: "", } if err := node.Insert(); err != nil { t.Fatalf("failed to seed relay node: %v", err) @@ -25,8 +25,8 @@ func TestHeartbeatRelayPersistsRuntimeAndObservability(t *testing.T) { now := time.Now().UTC() _, err := HeartbeatRelay(node, RelayHeartbeatPayload{ - RelayVersion: "v0.1.0", - FrpVersion: "0.61.0", + Version: "v0.1.0", + ExtVersion: "0.61.0", RelayStatus: "healthy", FrpsConnCount: 7, FrpsProxyCount: 3, @@ -62,11 +62,8 @@ func TestHeartbeatRelayPersistsRuntimeAndObservability(t *testing.T) { 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) + if updated.Version != "v0.1.0" || updated.ExtVersion != "0.61.0" { + t.Fatalf("expected relay versions to be updated, got relay=%q frp=%q", updated.Version, updated.ExtVersion) } profile, err := model.GetNodeSystemProfile(node.NodeID) 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 2ee728e5..eb1d5d63 100644 --- a/openflare_server/web/features/nodes/components/node-detail-page.tsx +++ b/openflare_server/web/features/nodes/components/node-detail-page.tsx @@ -1239,8 +1239,8 @@ function EdgeNodeDetailPage({ node }: { node: NodeItem }) {
-

Agent:{node.agent_version || 'unknown'}

-

Nginx:{node.nginx_version || 'unknown'}

+

Agent:{node.version || 'unknown'}

+

Nginx:{node.ext_version || 'unknown'}

当前配置:{node.current_version || '未应用'}

@@ -1742,7 +1742,7 @@ function EdgeNodeDetailPage({ node }: { node: NodeItem }) {

- {node.agent_version || 'unknown'} + {node.version || 'unknown'}

diff --git a/openflare_server/web/features/nodes/components/nodes-page.tsx b/openflare_server/web/features/nodes/components/nodes-page.tsx index 0f0a2c3e..55b5d9d7 100644 --- a/openflare_server/web/features/nodes/components/nodes-page.tsx +++ b/openflare_server/web/features/nodes/components/nodes-page.tsx @@ -344,9 +344,7 @@ export function NodesPage() { /> - {node.node_type === 'tunnel_relay' - ? node.relay_version || node.agent_version || 'unknown' - : node.agent_version || 'unknown'} + {node.version || 'unknown'}
diff --git a/openflare_server/web/features/nodes/components/relay-detail-page.tsx b/openflare_server/web/features/nodes/components/relay-detail-page.tsx index b511dbfa..2b72ec11 100644 --- a/openflare_server/web/features/nodes/components/relay-detail-page.tsx +++ b/openflare_server/web/features/nodes/components/relay-detail-page.tsx @@ -47,8 +47,6 @@ import { buildRelayDockerInstallCommand, getApplyLabel, getApplyVariant, - getNodeStatusLabel, - getNodeStatusVariant, getServerUrl, getUpdateMode, isMeaningfulTime, @@ -839,8 +837,8 @@ export function RelayDetailPage({ node }: { node: NodeItem }) {
-

Relay 中继版本:{node.relay_version || 'unknown'}

-

frps 核心版本:{node.relay_frp_version || 'unknown'}

+

Relay 中继版本:{node.version || 'unknown'}

+

frps 核心版本:{node.ext_version || 'unknown'}

中继网络接入:{node.relay_agent_access_addr || '—'}

@@ -1266,7 +1264,7 @@ export function RelayDetailPage({ node }: { node: NodeItem }) {

- {node.agent_version || 'unknown'} + {node.version || 'unknown'}

diff --git a/openflare_server/web/features/nodes/components/tunnel-detail-page.tsx b/openflare_server/web/features/nodes/components/tunnel-detail-page.tsx index 4fec24d4..9aa73b71 100644 --- a/openflare_server/web/features/nodes/components/tunnel-detail-page.tsx +++ b/openflare_server/web/features/nodes/components/tunnel-detail-page.tsx @@ -771,7 +771,7 @@ export function TunnelDetailPage({ node }: { node: NodeItem }) {
-

Client (openflared): {node.agent_version || 'unknown'}

+

Client (openflared): {node.version || 'unknown'}

当前配置版本:{node.current_version || '未应用'}

节点类型:隧道客户端 (tunnel_client)

@@ -1192,7 +1192,7 @@ export function TunnelDetailPage({ node }: { node: NodeItem }) {

- {node.agent_version || 'unknown'} + {node.version || 'unknown'}

diff --git a/openflare_server/web/features/nodes/types.ts b/openflare_server/web/features/nodes/types.ts index 974c484f..97d800f5 100644 --- a/openflare_server/web/features/nodes/types.ts +++ b/openflare_server/web/features/nodes/types.ts @@ -13,8 +13,6 @@ 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; @@ -27,8 +25,8 @@ export interface NodeItem { update_channel: ReleaseChannel; update_tag: string; restart_openresty_requested: boolean; - agent_version: string; - nginx_version: string; + version: string; + ext_version: string; openresty_status: 'healthy' | 'unhealthy' | 'unknown'; openresty_message: string; status: 'online' | 'offline' | 'pending';