[功能] 添加迁移遗留观察性列的功能,支持从 raw_json 填充 metadata_json

This commit is contained in:
ryan
2026-03-19 16:31:05 +08:00
parent dd49b2777d
commit f26fcd028e
10 changed files with 124 additions and 22 deletions
+48
View File
@@ -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
}
+48
View File
@@ -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")
}
}
@@ -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"`
}
+1 -1
View File
@@ -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"`
}
@@ -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"`
}
@@ -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"`
}
@@ -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
@@ -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) {
+2 -6
View File
@@ -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
+11 -10
View File
@@ -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;