From f26fcd028eee7255ff8e6aecdbd56ddd6b204ef9 Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 19 Mar 2026 16:31:05 +0800 Subject: [PATCH] =?UTF-8?q?[=E5=8A=9F=E8=83=BD]=20=E6=B7=BB=E5=8A=A0?= =?UTF-8?q?=E8=BF=81=E7=A7=BB=E9=81=97=E7=95=99=E8=A7=82=E5=AF=9F=E6=80=A7?= =?UTF-8?q?=E5=88=97=E7=9A=84=E5=8A=9F=E8=83=BD=EF=BC=8C=E6=94=AF=E6=8C=81?= =?UTF-8?q?=E4=BB=8E=20raw=5Fjson=20=E5=A1=AB=E5=85=85=20metadata=5Fjson?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- openflare_server/model/main.go | 48 +++++++++++++++++++ openflare_server/model/main_test.go | 48 +++++++++++++++++++ openflare_server/model/node_access_log.go | 1 - openflare_server/model/node_health_event.go | 2 +- .../model/node_metric_snapshot.go | 1 - openflare_server/model/node_request_report.go | 1 - openflare_server/model/node_system_profile.go | 2 - openflare_server/service/node_update_test.go | 14 ++++++ openflare_server/service/observability.go | 8 +--- openflare_server/web/features/nodes/types.ts | 21 ++++---- 10 files changed, 124 insertions(+), 22 deletions(-) diff --git a/openflare_server/model/main.go b/openflare_server/model/main.go index d2308814..464d26e2 100644 --- a/openflare_server/model/main.go +++ b/openflare_server/model/main.go @@ -1,6 +1,7 @@ package model import ( + "encoding/json" "fmt" "github.com/glebarez/sqlite" "gorm.io/driver/postgres" @@ -149,6 +150,50 @@ func migrateTextColumns(db *gorm.DB, backend string) error { return nil } +func migrateObservabilityLegacyColumns(db *gorm.DB) error { + if db == nil { + return nil + } + if !db.Migrator().HasTable(&NodeHealthEvent{}) || !db.Migrator().HasColumn(&NodeHealthEvent{}, "raw_json") { + return nil + } + type legacyHealthEventRaw struct { + ID uint + RawJSON string + MetadataJSON string + } + type legacyHealthEventPayload struct { + Metadata map[string]string `json:"metadata"` + } + + var rows []legacyHealthEventRaw + if err := db.Model(&NodeHealthEvent{}). + Select("id, raw_json, metadata_json"). + Where("raw_json <> '' AND (metadata_json IS NULL OR metadata_json = '')"). + Find(&rows).Error; err != nil { + return fmt.Errorf("query legacy node health event raw_json failed: %w", err) + } + for _, row := range rows { + var payload legacyHealthEventPayload + if err := json.Unmarshal([]byte(row.RawJSON), &payload); err != nil { + continue + } + if len(payload.Metadata) == 0 { + continue + } + metadataJSON, err := json.Marshal(payload.Metadata) + if err != nil { + continue + } + if err := db.Model(&NodeHealthEvent{}). + Where("id = ?", row.ID). + Update("metadata_json", string(metadataJSON)).Error; err != nil { + return fmt.Errorf("migrate node health event metadata_json failed: %w", err) + } + } + return nil +} + func isDatabaseEmpty(db *gorm.DB) (bool, error) { for _, item := range registeredModels() { var count int64 @@ -310,6 +355,9 @@ func InitDB() (err error) { if err = migrateTextColumns(db, backend); err != nil { return err } + if err = migrateObservabilityLegacyColumns(db); err != nil { + return err + } if err = migrateSQLiteDataIfNeeded(db, backend); err != nil { return err } diff --git a/openflare_server/model/main_test.go b/openflare_server/model/main_test.go index 396f934b..5c6664c3 100644 --- a/openflare_server/model/main_test.go +++ b/openflare_server/model/main_test.go @@ -1,8 +1,10 @@ package model import ( + "encoding/json" "path/filepath" "testing" + "time" "github.com/glebarez/sqlite" "gorm.io/gorm" @@ -154,3 +156,49 @@ func TestRegisterShardingAutoMigratesShardTables(t *testing.T) { } } } + +func TestMigrateObservabilityLegacyColumnsBackfillsHealthEventMetadata(t *testing.T) { + db := openTestSQLiteDB(t, "legacy-health-events.db") + + if err := db.Exec("ALTER TABLE node_health_events ADD COLUMN raw_json TEXT").Error; err != nil { + t.Fatalf("add raw_json column: %v", err) + } + rawJSON, err := json.Marshal(map[string]any{ + "event_type": "sync_error", + "metadata": map[string]string{ + "reason": "checksum_mismatch", + "scope": "routes", + }, + }) + if err != nil { + t.Fatalf("marshal raw json: %v", err) + } + event := &NodeHealthEvent{ + NodeID: "node-legacy", + EventType: "sync_error", + Severity: "warning", + Status: "active", + Message: "checksum mismatch", + FirstTriggeredAt: time.Now().Add(-time.Minute), + LastTriggeredAt: time.Now(), + ReportedAt: time.Now(), + } + if err := db.Create(event).Error; err != nil { + t.Fatalf("create health event: %v", err) + } + if err := db.Exec("UPDATE node_health_events SET raw_json = ? WHERE id = ?", string(rawJSON), event.ID).Error; err != nil { + t.Fatalf("seed legacy raw_json: %v", err) + } + + if err := migrateObservabilityLegacyColumns(db); err != nil { + t.Fatalf("migrateObservabilityLegacyColumns: %v", err) + } + + var got NodeHealthEvent + if err := db.First(&got, event.ID).Error; err != nil { + t.Fatalf("query health event: %v", err) + } + if got.MetadataJSON == "" { + t.Fatal("expected metadata_json to be backfilled") + } +} diff --git a/openflare_server/model/node_access_log.go b/openflare_server/model/node_access_log.go index 9853a45a..480e3e27 100644 --- a/openflare_server/model/node_access_log.go +++ b/openflare_server/model/node_access_log.go @@ -18,7 +18,6 @@ type NodeAccessLog struct { Host string `json:"host" gorm:"index;size:255"` Path string `json:"path" gorm:"size:2048"` StatusCode int `json:"status_code" gorm:"index"` - RawJSON string `json:"raw_json" gorm:"type:text"` CreatedAt time.Time `json:"created_at"` } diff --git a/openflare_server/model/node_health_event.go b/openflare_server/model/node_health_event.go index 4bc0df63..88fdc0f7 100644 --- a/openflare_server/model/node_health_event.go +++ b/openflare_server/model/node_health_event.go @@ -13,7 +13,7 @@ type NodeHealthEvent struct { LastTriggeredAt time.Time `json:"last_triggered_at" gorm:"index"` ReportedAt time.Time `json:"reported_at" gorm:"index"` ResolvedAt *time.Time `json:"resolved_at" gorm:"index"` - RawJSON string `json:"raw_json" gorm:"type:text"` + MetadataJSON string `json:"metadata_json" gorm:"type:text"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` } diff --git a/openflare_server/model/node_metric_snapshot.go b/openflare_server/model/node_metric_snapshot.go index 1e956a12..7f3aac48 100644 --- a/openflare_server/model/node_metric_snapshot.go +++ b/openflare_server/model/node_metric_snapshot.go @@ -23,7 +23,6 @@ type NodeMetricSnapshot struct { OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` OpenrestyConnections int64 `json:"openresty_connections"` - RawJSON string `json:"raw_json" gorm:"type:text"` CreatedAt time.Time `json:"created_at"` } diff --git a/openflare_server/model/node_request_report.go b/openflare_server/model/node_request_report.go index 8fe5f2aa..fb5d17bf 100644 --- a/openflare_server/model/node_request_report.go +++ b/openflare_server/model/node_request_report.go @@ -18,7 +18,6 @@ type NodeRequestReport struct { StatusCodesJSON string `json:"status_codes_json" gorm:"type:text"` TopDomainsJSON string `json:"top_domains_json" gorm:"type:text"` SourceCountriesJSON string `json:"source_countries_json" gorm:"type:text"` - RawJSON string `json:"raw_json" gorm:"type:text"` CreatedAt time.Time `json:"created_at"` } diff --git a/openflare_server/model/node_system_profile.go b/openflare_server/model/node_system_profile.go index 01866aeb..5f2b345f 100644 --- a/openflare_server/model/node_system_profile.go +++ b/openflare_server/model/node_system_profile.go @@ -20,7 +20,6 @@ type NodeSystemProfile struct { TotalDiskBytes int64 `json:"total_disk_bytes"` UptimeSeconds int64 `json:"uptime_seconds"` ReportedAt time.Time `json:"reported_at" gorm:"index"` - RawJSON string `json:"raw_json" gorm:"type:text"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` } @@ -49,7 +48,6 @@ func UpsertNodeSystemProfile(profile *NodeSystemProfile) error { "total_disk_bytes", "uptime_seconds", "reported_at", - "raw_json", "updated_at", }), }).Create(profile).Error diff --git a/openflare_server/service/node_update_test.go b/openflare_server/service/node_update_test.go index 3d5d1e58..a839aa9c 100644 --- a/openflare_server/service/node_update_test.go +++ b/openflare_server/service/node_update_test.go @@ -1,6 +1,7 @@ package service import ( + "encoding/json" "io" "net" "net/http" @@ -677,6 +678,9 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) { Severity: NodeHealthSeverityCritical, Message: "reload failed", TriggeredAtUnix: time.Now().Add(-2 * time.Minute).Unix(), + Metadata: map[string]string{ + "source": "runtime", + }, }, }, }) @@ -731,6 +735,16 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) { if len(events) != 1 || events[0].EventType != "openresty_unhealthy" { t.Fatalf("unexpected active health events: %+v", events) } + if events[0].MetadataJSON == "" { + t.Fatal("expected metadata_json to persist") + } + var metadata map[string]string + if err := json.Unmarshal([]byte(events[0].MetadataJSON), &metadata); err != nil { + t.Fatalf("expected metadata_json to be valid json: %v", err) + } + if metadata["source"] != "runtime" { + t.Fatalf("unexpected metadata json: %+v", metadata) + } } func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) { diff --git a/openflare_server/service/observability.go b/openflare_server/service/observability.go index 84cb10ff..48544bf1 100644 --- a/openflare_server/service/observability.go +++ b/openflare_server/service/observability.go @@ -152,7 +152,6 @@ func persistNodeSystemProfile(tx *gorm.DB, nodeID string, profile *AgentNodeSyst TotalDiskBytes: profile.TotalDiskBytes, UptimeSeconds: profile.UptimeSeconds, ReportedAt: timeFromUnix(profile.ReportedAtUnix, reportedAt), - RawJSON: marshalJSON(profile), } return tx.Model(&model.NodeSystemProfile{}).Where("node_id = ?", nodeID).Assign(record).FirstOrCreate(record).Error } @@ -176,7 +175,6 @@ func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *AgentNodeMe OpenrestyRxBytes: snapshot.OpenrestyRxBytes, OpenrestyTxBytes: snapshot.OpenrestyTxBytes, OpenrestyConnections: snapshot.OpenrestyConnections, - RawJSON: marshalJSON(snapshot), } return tx.Where("node_id = ? AND captured_at = ?", nodeID, record.CapturedAt).FirstOrCreate(record).Error } @@ -198,7 +196,6 @@ func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *AgentNodeTraff StatusCodesJSON: marshalJSON(report.StatusCodes), TopDomainsJSON: marshalJSON(report.TopDomains), SourceCountriesJSON: marshalJSON(report.SourceCountries), - RawJSON: marshalJSON(report), } return tx.Where("node_id = ? AND window_started_at = ? AND window_ended_at = ?", nodeID, record.WindowStartedAt, record.WindowEndedAt).FirstOrCreate(record).Error } @@ -223,7 +220,6 @@ func persistNodeAccessLogs(tx *gorm.DB, nodeID string, logs []AgentNodeAccessLog Host: strings.TrimSpace(item.Host), Path: truncateForDatabase(strings.TrimSpace(item.Path), nodeAccessLogPathMaxLength), StatusCode: item.StatusCode, - RawJSON: marshalJSON(item), } if resolver != nil { record.Region = resolver.Resolve(record.RemoteAddr) @@ -275,7 +271,7 @@ func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHea existing.Message = normalizeHealthEventMessage(event.Message) existing.LastTriggeredAt = triggeredAt existing.ReportedAt = reportedAt - existing.RawJSON = marshalJSON(event) + existing.MetadataJSON = marshalJSON(event.Metadata) existing.ResolvedAt = nil if err := tx.Save(existing).Error; err != nil { return err @@ -291,7 +287,7 @@ func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHea FirstTriggeredAt: triggeredAt, LastTriggeredAt: triggeredAt, ReportedAt: reportedAt, - RawJSON: marshalJSON(event), + MetadataJSON: marshalJSON(event.Metadata), } if err := tx.Create(record).Error; err != nil { return err diff --git a/openflare_server/web/features/nodes/types.ts b/openflare_server/web/features/nodes/types.ts index 1555f6b0..f7ae80df 100644 --- a/openflare_server/web/features/nodes/types.ts +++ b/openflare_server/web/features/nodes/types.ts @@ -183,16 +183,17 @@ export interface NodeObservabilityTrends { disk_io_24h: NodeDiskIOTrendPoint[]; } -export interface NodeHealthEvent { - event_type: string; - severity: string; - status: string; - message: string; - first_triggered_at: string; - last_triggered_at: string; - reported_at: string; - resolved_at?: string | null; -} +export interface NodeHealthEvent { + event_type: string; + severity: string; + status: string; + message: string; + metadata_json?: string; + first_triggered_at: string; + last_triggered_at: string; + reported_at: string; + resolved_at?: string | null; +} export interface NodeObservability { node_id: string;