From b968043117261014fd9f0c20b330c102d06d6cfe Mon Sep 17 00:00:00 2001 From: ryan Date: Mon, 1 Jun 2026 16:19:59 +0800 Subject: [PATCH] =?UTF-8?q?[=E4=BC=98=E5=8C=96]=20=E4=BF=AE=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- openflare_server/service/observability.go | 25 +++++++++++++++++- openflare_server/service/relay.go | 31 +++++++++++++++++++++++ 2 files changed, 55 insertions(+), 1 deletion(-) diff --git a/openflare_server/service/observability.go b/openflare_server/service/observability.go index 0f6dc2d7..a434cdd1 100644 --- a/openflare_server/service/observability.go +++ b/openflare_server/service/observability.go @@ -276,12 +276,21 @@ func persistNodeAccessLogs(tx *gorm.DB, nodeID string, logs []AgentNodeAccessLog } func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHealthEvent, reportedAt time.Time) error { + return reconcileScopedNodeHealthEvents(tx, nodeID, events, reportedAt, nil) +} + +func reconcileScopedNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHealthEvent, reportedAt time.Time, managedEventTypes map[string]struct{}) error { activeTypes := make(map[string]AgentNodeHealthEvent, len(events)) for _, event := range events { eventType := normalizeHealthEventType(event.EventType) if eventType == "" { continue } + if len(managedEventTypes) > 0 { + if _, ok := managedEventTypes[eventType]; !ok { + continue + } + } event.EventType = eventType event.Severity = normalizeHealthSeverity(event.Severity) if event.TriggeredAtUnix <= 0 { @@ -291,7 +300,21 @@ func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHea } var activeEvents []*model.NodeHealthEvent - if err := tx.Where("node_id = ? AND status = ?", nodeID, NodeHealthEventStatusActive).Find(&activeEvents).Error; err != nil { + query := tx.Where("node_id = ? AND status = ?", nodeID, NodeHealthEventStatusActive) + if len(managedEventTypes) > 0 { + scopedTypes := make([]string, 0, len(managedEventTypes)) + for eventType := range managedEventTypes { + eventType = normalizeHealthEventType(eventType) + if eventType != "" { + scopedTypes = append(scopedTypes, eventType) + } + } + if len(scopedTypes) == 0 { + return nil + } + query = query.Where("event_type IN ?", scopedTypes) + } + if err := query.Find(&activeEvents).Error; err != nil { return err } diff --git a/openflare_server/service/relay.go b/openflare_server/service/relay.go index 153a6a20..18f343ed 100644 --- a/openflare_server/service/relay.go +++ b/openflare_server/service/relay.go @@ -7,6 +7,8 @@ import ( "openflare/model" "strings" "time" + + "gorm.io/gorm" ) // RelayHeartbeatPayload is the payload sent by OpenFlareRelay in each heartbeat. @@ -23,6 +25,8 @@ type RelayHeartbeatPayload struct { HealthEvents []AgentNodeHealthEvent `json:"health_events,omitempty"` } +const relayFrpsUnhealthyEventType = "frps_unhealthy" + // RelayConfig is the frps configuration sent to the Relay. type RelayConfig struct { BindPort int `json:"bind_port"` @@ -98,6 +102,9 @@ func HeartbeatRelay(node *model.Node, payload RelayHeartbeatPayload) (*RelayHear return nil, fmt.Errorf("update relay heartbeat: %w", err) } } + if err := reconcileRelayHealthEvents(node.NodeID, payload.RelayStatus, now); err != nil { + return nil, fmt.Errorf("reconcile relay health events: %w", err) + } refreshAccessTokenCache(node) persistRelayHeartbeatObservability(node.NodeID, payload, node.LastSeenAt) @@ -142,6 +149,30 @@ func buildRelaySettings() *RelaySettings { } } +func reconcileRelayHealthEvents(nodeID string, relayStatus string, reportedAt time.Time) error { + if relayStatus == "unknown" { + return nil + } + managedTypes := map[string]struct{}{ + relayFrpsUnhealthyEventType: {}, + } + events := []AgentNodeHealthEvent{} + if relayStatus == "unhealthy" { + events = append(events, AgentNodeHealthEvent{ + EventType: relayFrpsUnhealthyEventType, + Severity: NodeHealthSeverityCritical, + Message: "frps runtime is not healthy", + TriggeredAtUnix: reportedAt.Unix(), + Metadata: map[string]string{ + "relay_status": relayStatus, + }, + }) + } + return model.DB.Transaction(func(tx *gorm.DB) error { + return reconcileScopedNodeHealthEvents(tx, nodeID, events, reportedAt, managedTypes) + }) +} + func normalizeRelayStatus(status string) string { switch strings.ToLower(strings.TrimSpace(status)) { case "healthy":