From 5943372ce43e38e91a52b5f75ed27786223d0062 Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 18 Jun 2026 16:38:24 +0800 Subject: [PATCH] migrate db --- Wavelet/.env.example | 8 +- Wavelet/config.example.yaml | 6 +- Wavelet/docker-compose.yml | 6 +- .../internal/apps/admin/db_manage/routers.go | 2 +- Wavelet/internal/apps/admin/status/routers.go | 10 +- .../internal/apps/openflare/agent/logics.go | 6 + .../apps/openflare/agent/observability.go | 450 ++++++++++++++++++ .../internal/apps/openflare/agent/types.go | 27 +- ...6190010_create_of_observability_tables.sql | 136 ++++++ ...6190010_create_of_observability_tables.sql | 136 ++++++ Wavelet/internal/db/postgres.go | 2 +- .../sql/create_clickhouse_risk.sql | 4 +- docker-compose.yaml | 2 +- docs/changelog/index.md | 3 + 14 files changed, 768 insertions(+), 30 deletions(-) create mode 100644 Wavelet/internal/apps/openflare/agent/observability.go create mode 100644 Wavelet/internal/db/migrator/goose/postgres/202606190010_create_of_observability_tables.sql create mode 100644 Wavelet/internal/db/migrator/goose/sqlite/202606190010_create_of_observability_tables.sql diff --git a/Wavelet/.env.example b/Wavelet/.env.example index e9af3f65..92d94e6b 100644 --- a/Wavelet/.env.example +++ b/Wavelet/.env.example @@ -13,7 +13,7 @@ JAEGER_OTLP_GRPC_PORT=4317 JAEGER_OTLP_HTTP_PORT=4318 # ─── PostgreSQL 容器配置(仅 docker-compose 使用)──────────────────────────── -POSTGRES_DB=wavelet +POSTGRES_DB=openflare POSTGRES_USER=postgres POSTGRES_PASSWORD=postgres @@ -39,12 +39,12 @@ APP_SESSION_SECURE=true # 设置 DB_HOST 后自动启用 PostgreSQL,也可通过 DB_ENABLED 显式控制 # DB_ENABLED=false 时使用 SQLite 作为后备数据库 DB_ENABLED=true -# SQLITE_PATH=./data/wavelet.db +# SQLITE_PATH=./data/openflare.db DB_HOST=postgres DB_PORT=5432 DB_USERNAME=postgres DB_PASSWORD=postgres -DB_NAME=wavelet +DB_NAME=openflare DB_SSL_MODE=disable DB_TIMEZONE=Asia/Shanghai # DB_LOG_LEVEL=info @@ -67,7 +67,7 @@ REDIS_KEY_PREFIX=wavelet: # CLICKHOUSE_HOST=clickhouse:9000 # CLICKHOUSE_USERNAME=default # CLICKHOUSE_PASSWORD= -# CLICKHOUSE_NAME=wavelet +# CLICKHOUSE_NAME=openflare # ─── 日志 ────────────────────────────────────────────────────────────────────── LOG_LEVEL=info diff --git a/Wavelet/config.example.yaml b/Wavelet/config.example.yaml index 090c2745..b2a1d41e 100644 --- a/Wavelet/config.example.yaml +++ b/Wavelet/config.example.yaml @@ -21,12 +21,12 @@ app: # Supports Standalone and Primary-Replica (read/write split) modes. database: enabled: true - sqlite_path: "wavelet.db" # PostgreSQL 禁用时使用此 SQLite 文件路径 + sqlite_path: "openflare.db" # PostgreSQL 禁用时使用此 SQLite 文件路径 host: "127.0.0.1" port: 5432 username: "postgres" password: "postgres" - database: "wavelet" + database: "openflare" max_idle_conn: 16 max_open_conn: 128 conn_max_lifetime: 1800 @@ -105,7 +105,7 @@ clickhouse: - "127.0.0.1:9000" username: "default" password: "" - database: "wavelet" + database: "openflare" max_idle_conn: 10 max_open_conn: 100 conn_max_lifetime: 3600 diff --git a/Wavelet/docker-compose.yml b/Wavelet/docker-compose.yml index 82606f2e..4859467e 100644 --- a/Wavelet/docker-compose.yml +++ b/Wavelet/docker-compose.yml @@ -30,7 +30,7 @@ services: image: postgres:17-alpine restart: unless-stopped environment: - POSTGRES_DB: ${POSTGRES_DB:-wavelet} + POSTGRES_DB: ${POSTGRES_DB:-openflare} POSTGRES_USER: ${POSTGRES_USER:-postgres} POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:-postgres} TZ: ${TZ:-Asia/Shanghai} @@ -39,7 +39,7 @@ services: volumes: - ./data/postgres_data:/var/lib/postgresql/data healthcheck: - test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-postgres} -d ${POSTGRES_DB:-wavelet}"] + test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-postgres} -d ${POSTGRES_DB:-openflare}"] interval: 10s timeout: 5s retries: 5 @@ -76,7 +76,7 @@ services: # profiles: # - clickhouse # environment: -# CLICKHOUSE_DB: ${CLICKHOUSE_DB:-wavelet} +# CLICKHOUSE_DB: ${CLICKHOUSE_DB:-openflare} # CLICKHOUSE_USER: ${CLICKHOUSE_USER:-default} # CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-123456} # CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: 1 diff --git a/Wavelet/internal/apps/admin/db_manage/routers.go b/Wavelet/internal/apps/admin/db_manage/routers.go index e2a1fbc0..3f7086c9 100644 --- a/Wavelet/internal/apps/admin/db_manage/routers.go +++ b/Wavelet/internal/apps/admin/db_manage/routers.go @@ -106,7 +106,7 @@ func formatBytes(bytes uint64) string { func getSQLiteOverview(gormDB *gorm.DB) (DBOverviewResponse, error) { name := config.Config.Database.SQLitePath if name == "" { - name = "./data/wavelet.db" + name = "./data/openflare.db" } var version string diff --git a/Wavelet/internal/apps/admin/status/routers.go b/Wavelet/internal/apps/admin/status/routers.go index 29a1f23f..633319df 100644 --- a/Wavelet/internal/apps/admin/status/routers.go +++ b/Wavelet/internal/apps/admin/status/routers.go @@ -213,7 +213,7 @@ func getSQLiteInfo(ctx context.Context) DatabaseInfoResponse { Version: "SQLite", } if info.Name == "" { - info.Name = "./data/wavelet.db" + info.Name = "./data/openflare.db" } gormDB := db.DB(ctx) if gormDB == nil { @@ -287,7 +287,7 @@ func ExportDatabase(c *gin.Context) { func exportSQLite(c *gin.Context) { path := config.Config.Database.SQLitePath if path == "" { - path = "./data/wavelet.db" + path = "./data/openflare.db" } f, err := os.Open(path) //nolint:gosec // path is loaded from server startup configuration, not user input @@ -307,11 +307,11 @@ func exportSQLite(c *gin.Context) { return } - c.Header("Content-Disposition", `attachment; filename="wavelet.db"`) + c.Header("Content-Disposition", `attachment; filename="openflare.db"`) c.Header("Content-Type", "application/octet-stream") c.Header("Content-Length", fmt.Sprintf("%d", fi.Size())) c.Status(http.StatusOK) - http.ServeContent(c.Writer, c.Request, "wavelet.db", fi.ModTime(), f) + http.ServeContent(c.Writer, c.Request, "openflare.db", fi.ModTime(), f) } // exportPostgres 执行 pg_dump 并将输出流式传输给客户端 @@ -340,7 +340,7 @@ func exportPostgres(c *gin.Context) { cmd.Env = os.Environ() } - fileName := fmt.Sprintf("wavelet_%s.sql", time.Now().Format("20060102_150405")) + fileName := fmt.Sprintf("openflare_%s.sql", time.Now().Format("20060102_150405")) c.Header("Content-Disposition", `attachment; filename="`+fileName+`"`) c.Header("Content-Type", "application/octet-stream") c.Status(http.StatusOK) diff --git a/Wavelet/internal/apps/openflare/agent/logics.go b/Wavelet/internal/apps/openflare/agent/logics.go index 843d08df..fefd0eee 100644 --- a/Wavelet/internal/apps/openflare/agent/logics.go +++ b/Wavelet/internal/apps/openflare/agent/logics.go @@ -118,6 +118,12 @@ func HeartbeatNode(ctx context.Context, authNode *model.OpenFlareNode, payload N refreshAccessTokenCache(ctx, authNode) + reportedAt := time.Now() + if authNode.LastSeenAt != nil { + reportedAt = *authNode.LastSeenAt + } + persistHeartbeatObservability(ctx, authNode.NodeID, payload, reportedAt) + activeConfig, err := getActiveConfigMeta(ctx) if err != nil && !isActiveConfigNotFound(err) { return nil, err diff --git a/Wavelet/internal/apps/openflare/agent/observability.go b/Wavelet/internal/apps/openflare/agent/observability.go new file mode 100644 index 00000000..0ef76935 --- /dev/null +++ b/Wavelet/internal/apps/openflare/agent/observability.go @@ -0,0 +1,450 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package agent + +import ( + "context" + "encoding/json" + "errors" + "strings" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "go.uber.org/zap" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +const ( + healthEventStatusActive = "active" + healthEventStatusResolved = "resolved" + healthSeverityInfo = "info" + healthSeverityWarning = "warning" + healthSeverityCritical = "critical" + accessLogPathMaxLength = 100 +) + +// NodeSystemProfile is the agent-reported system profile. +type NodeSystemProfile struct { + Hostname string `json:"hostname"` + OSName string `json:"os_name"` + OSVersion string `json:"os_version"` + KernelVersion string `json:"kernel_version"` + Architecture string `json:"architecture"` + CPUModel string `json:"cpu_model"` + CPUCores int `json:"cpu_cores"` + TotalMemoryBytes int64 `json:"total_memory_bytes"` + TotalDiskBytes int64 `json:"total_disk_bytes"` + UptimeSeconds int64 `json:"uptime_seconds"` + ReportedAtUnix int64 `json:"reported_at_unix"` +} + +// NodeMetricSnapshot is the agent-reported capacity snapshot. +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"` +} + +// NodeOpenrestyObservation is the agent-reported openresty network observation. +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"` +} + +// NodeTrafficReport is the agent-reported traffic window. +type NodeTrafficReport struct { + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + WindowEndedAtUnix int64 `json:"window_ended_at_unix"` + RequestCount int64 `json:"request_count"` + ErrorCount int64 `json:"error_count"` + UniqueVisitorCount int64 `json:"unique_visitor_count"` + StatusCodes map[string]int64 `json:"status_codes"` + TopDomains map[string]int64 `json:"top_domains"` + SourceCountries map[string]int64 `json:"source_countries"` +} + +// NodeAccessLog is a single access log row from the agent. +type NodeAccessLog struct { + LoggedAtUnix int64 `json:"logged_at_unix"` + RemoteAddr string `json:"remote_addr"` + Host string `json:"host"` + Path string `json:"path"` + StatusCode int `json:"status_code"` +} + +// BufferedObservabilityRecord is a buffered observability window from the agent. +type BufferedObservabilityRecord struct { + 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"` +} + +// NodeHealthEvent is an agent-reported health event. +type NodeHealthEvent struct { + EventType string `json:"event_type"` + Severity string `json:"severity"` + Message string `json:"message"` + TriggeredAtUnix int64 `json:"triggered_at_unix"` + Metadata map[string]string `json:"metadata"` +} + +func persistHeartbeatObservability(ctx context.Context, nodeID string, payload NodePayload, reportedAt time.Time) { + if strings.TrimSpace(nodeID) == "" { + return + } + if payload.Profile == nil && + payload.Snapshot == nil && + payload.TrafficReport == nil && + len(payload.AccessLogs) == 0 && + len(payload.BufferedObservability) == 0 && + payload.HealthEvents == nil { + return + } + + conn := db.DB(ctx) + if conn == nil { + return + } + + if err := conn.Transaction(func(tx *gorm.DB) error { + if err := persistNodeSystemProfile(tx, nodeID, payload.Profile, reportedAt); err != nil { + return err + } + if err := persistBufferedObservability(tx, nodeID, payload.BufferedObservability, reportedAt); err != nil { + return err + } + 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 + } + if err := persistNodeAccessLogs(tx, nodeID, payload.AccessLogs, reportedAt); err != nil { + return err + } + if payload.HealthEvents != nil { + if err := reconcileNodeHealthEvents(tx, nodeID, payload.HealthEvents, reportedAt); err != nil { + return err + } + } + return nil + }); err != nil { + zap.L().Error("persist heartbeat observability failed", zap.String("node_id", nodeID), zap.Error(err)) + } +} + +func persistBufferedObservability(tx *gorm.DB, nodeID string, records []BufferedObservabilityRecord, reportedAt time.Time) error { + for _, record := range records { + 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 + } + if err := persistNodeAccessLogs(tx, nodeID, record.AccessLogs, reportedAt); err != nil { + return err + } + } + return nil +} + +func persistNodeSystemProfile(tx *gorm.DB, nodeID string, profile *NodeSystemProfile, reportedAt time.Time) error { + if profile == nil { + return nil + } + record := &model.OpenFlareNodeSystemProfile{ + NodeID: nodeID, + Hostname: strings.TrimSpace(profile.Hostname), + OSName: strings.TrimSpace(profile.OSName), + OSVersion: strings.TrimSpace(profile.OSVersion), + KernelVersion: strings.TrimSpace(profile.KernelVersion), + Architecture: strings.TrimSpace(profile.Architecture), + CPUModel: strings.TrimSpace(profile.CPUModel), + CPUCores: profile.CPUCores, + TotalMemoryBytes: profile.TotalMemoryBytes, + TotalDiskBytes: profile.TotalDiskBytes, + UptimeSeconds: profile.UptimeSeconds, + ReportedAt: timeFromUnix(profile.ReportedAtUnix, reportedAt), + } + return tx.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "node_id"}}, + DoUpdates: clause.AssignmentColumns([]string{ + "hostname", + "os_name", + "os_version", + "kernel_version", + "architecture", + "cpu_model", + "cpu_cores", + "total_memory_bytes", + "total_disk_bytes", + "uptime_seconds", + "reported_at", + "updated_at", + }), + }).Create(record).Error +} + +func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *NodeMetricSnapshot, reportedAt time.Time) error { + if snapshot == nil { + return nil + } + record := &model.OpenFlareMetricSnapshot{ + 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 := metricSnapshotExists(tx, nodeID, record.CapturedAt) + if err != nil { + return err + } + if exists { + return nil + } + return tx.Create(record).Error +} + +func persistNodeOpenrestyObservation(tx *gorm.DB, nodeID string, obs *NodeOpenrestyObservation, reportedAt time.Time) error { + if obs == nil { + return nil + } + record := &model.OpenFlareNodeObservationOpenresty{ + 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 *NodeTrafficReport, reportedAt time.Time) error { + if report == nil { + return nil + } + if report.WindowEndedAtUnix > 0 && report.WindowStartedAtUnix > report.WindowEndedAtUnix { + return errors.New("traffic report window_started_at_unix 不能大于 window_ended_at_unix") + } + record := &model.OpenFlareRequestReport{ + NodeID: nodeID, + WindowStartedAt: timeFromUnix(report.WindowStartedAtUnix, reportedAt), + WindowEndedAt: timeFromUnix(report.WindowEndedAtUnix, reportedAt), + RequestCount: report.RequestCount, + ErrorCount: report.ErrorCount, + UniqueVisitorCount: report.UniqueVisitorCount, + StatusCodesJSON: marshalJSON(report.StatusCodes), + TopDomainsJSON: marshalJSON(report.TopDomains), + SourceCountriesJSON: marshalJSON(report.SourceCountries), + } + exists, err := requestReportExists(tx, nodeID, record.WindowStartedAt, record.WindowEndedAt) + if err != nil { + return err + } + if exists { + return nil + } + return tx.Create(record).Error +} + +func persistNodeAccessLogs(tx *gorm.DB, nodeID string, logs []NodeAccessLog, reportedAt time.Time) error { + for _, item := range logs { + record := &model.OpenFlareAccessLog{ + NodeID: nodeID, + LoggedAt: timeFromUnix(item.LoggedAtUnix, reportedAt), + RemoteAddr: strings.TrimSpace(item.RemoteAddr), + Host: strings.TrimSpace(item.Host), + Path: truncateForDatabase(strings.TrimSpace(item.Path), accessLogPathMaxLength), + StatusCode: item.StatusCode, + } + exists, err := accessLogExists(tx, record) + if err != nil { + return err + } + if exists { + continue + } + if err := tx.Create(record).Error; err != nil { + return err + } + } + return nil +} + +func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []NodeHealthEvent, reportedAt time.Time) error { + activeTypes := make(map[string]NodeHealthEvent, len(events)) + for _, event := range events { + eventType := normalizeHealthEventType(event.EventType) + if eventType == "" { + continue + } + event.EventType = eventType + event.Severity = normalizeHealthSeverity(event.Severity) + if event.TriggeredAtUnix <= 0 { + event.TriggeredAtUnix = reportedAt.Unix() + } + activeTypes[eventType] = event + } + + var activeEvents []*model.OpenFlareHealthEvent + if err := tx.Where("node_id = ? AND status = ?", nodeID, healthEventStatusActive).Find(&activeEvents).Error; err != nil { + return err + } + + activeByType := make(map[string]*model.OpenFlareHealthEvent, len(activeEvents)) + for _, event := range activeEvents { + activeByType[event.EventType] = event + } + + for eventType, event := range activeTypes { + triggeredAt := timeFromUnix(event.TriggeredAtUnix, reportedAt) + if existing, ok := activeByType[eventType]; ok { + existing.Severity = event.Severity + existing.Message = normalizeHealthEventMessage(event.Message) + existing.LastTriggeredAt = triggeredAt + existing.ReportedAt = reportedAt + existing.MetadataJSON = marshalJSON(event.Metadata) + existing.ResolvedAt = nil + if err := tx.Save(existing).Error; err != nil { + return err + } + continue + } + record := &model.OpenFlareHealthEvent{ + NodeID: nodeID, + EventType: eventType, + Severity: event.Severity, + Status: healthEventStatusActive, + Message: normalizeHealthEventMessage(event.Message), + FirstTriggeredAt: triggeredAt, + LastTriggeredAt: triggeredAt, + ReportedAt: reportedAt, + MetadataJSON: marshalJSON(event.Metadata), + } + if err := tx.Create(record).Error; err != nil { + return err + } + } + + for _, existing := range activeEvents { + if _, ok := activeTypes[existing.EventType]; ok { + continue + } + resolvedAt := reportedAt + existing.Status = healthEventStatusResolved + existing.ReportedAt = reportedAt + existing.ResolvedAt = &resolvedAt + if err := tx.Save(existing).Error; err != nil { + return err + } + } + + return nil +} + +func metricSnapshotExists(tx *gorm.DB, nodeID string, capturedAt time.Time) (bool, error) { + var count int64 + if err := tx.Model(&model.OpenFlareMetricSnapshot{}). + Where("node_id = ? AND captured_at = ?", nodeID, capturedAt). + Limit(1). + Count(&count).Error; err != nil { + return false, err + } + return count > 0, nil +} + +func requestReportExists(tx *gorm.DB, nodeID string, windowStartedAt, windowEndedAt time.Time) (bool, error) { + var count int64 + if err := tx.Model(&model.OpenFlareRequestReport{}). + Where("node_id = ? AND window_started_at = ? AND window_ended_at = ?", nodeID, windowStartedAt, windowEndedAt). + Limit(1). + Count(&count).Error; err != nil { + return false, err + } + return count > 0, nil +} + +func accessLogExists(tx *gorm.DB, record *model.OpenFlareAccessLog) (bool, error) { + var count int64 + if err := tx.Model(&model.OpenFlareAccessLog{}). + Where( + "node_id = ? AND logged_at = ? AND remote_addr = ? AND host = ? AND path = ? AND status_code = ?", + record.NodeID, + record.LoggedAt, + record.RemoteAddr, + record.Host, + record.Path, + record.StatusCode, + ). + Limit(1). + Count(&count).Error; err != nil { + return false, err + } + return count > 0, nil +} + +func normalizeHealthEventType(eventType string) string { + eventType = strings.TrimSpace(strings.ToLower(eventType)) + eventType = strings.ReplaceAll(eventType, " ", "_") + return eventType +} + +func normalizeHealthSeverity(severity string) string { + switch strings.ToLower(strings.TrimSpace(severity)) { + case healthSeverityCritical: + return healthSeverityCritical + case healthSeverityInfo: + return healthSeverityInfo + default: + return healthSeverityWarning + } +} + +func normalizeHealthEventMessage(message string) string { + return truncateForDatabase(message, 4096) +} + +func timeFromUnix(unixSeconds int64, fallback time.Time) time.Time { + if unixSeconds <= 0 { + return fallback + } + return time.Unix(unixSeconds, 0).UTC() +} + +func marshalJSON(value any) string { + if value == nil { + return "" + } + raw, err := json.Marshal(value) + if err != nil { + return "" + } + return string(raw) +} diff --git a/Wavelet/internal/apps/openflare/agent/types.go b/Wavelet/internal/apps/openflare/agent/types.go index e5a52f50..47d6e24a 100644 --- a/Wavelet/internal/apps/openflare/agent/types.go +++ b/Wavelet/internal/apps/openflare/agent/types.go @@ -18,16 +18,23 @@ const ( // NodePayload is the agent register/heartbeat payload. type NodePayload struct { - NodeID string `json:"node_id"` - Name string `json:"name"` - IP string `json:"ip"` - 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"` - WAFIPGroupChecksums map[string]string `json:"waf_ip_group_checksums,omitempty"` + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + 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 *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"` + HealthEvents []NodeHealthEvent `json:"health_events"` + WAFIPGroupChecksums map[string]string `json:"waf_ip_group_checksums,omitempty"` } // ApplyLogPayload is the agent apply log report payload. diff --git a/Wavelet/internal/db/migrator/goose/postgres/202606190010_create_of_observability_tables.sql b/Wavelet/internal/db/migrator/goose/postgres/202606190010_create_of_observability_tables.sql new file mode 100644 index 00000000..24d22a19 --- /dev/null +++ b/Wavelet/internal/db/migrator/goose/postgres/202606190010_create_of_observability_tables.sql @@ -0,0 +1,136 @@ +-- +goose Up +CREATE TABLE of_node_system_profiles ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + hostname VARCHAR(255) NOT NULL DEFAULT '', + os_name VARCHAR(128) NOT NULL DEFAULT '', + os_version VARCHAR(128) NOT NULL DEFAULT '', + kernel_version VARCHAR(128) NOT NULL DEFAULT '', + architecture VARCHAR(64) NOT NULL DEFAULT '', + cpu_model VARCHAR(255) NOT NULL DEFAULT '', + cpu_cores INTEGER NOT NULL DEFAULT 0, + total_memory_bytes BIGINT NOT NULL DEFAULT 0, + total_disk_bytes BIGINT NOT NULL DEFAULT 0, + uptime_seconds BIGINT NOT NULL DEFAULT 0, + reported_at TIMESTAMPTZ NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE UNIQUE INDEX idx_of_node_system_profiles_node_id ON of_node_system_profiles (node_id); +CREATE INDEX idx_of_node_system_profiles_reported_at ON of_node_system_profiles (reported_at); + +CREATE TABLE of_node_metric_snapshots ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + captured_at TIMESTAMPTZ NOT NULL, + cpu_usage_percent DOUBLE PRECISION NOT NULL DEFAULT 0, + memory_used_bytes BIGINT NOT NULL DEFAULT 0, + memory_total_bytes BIGINT NOT NULL DEFAULT 0, + storage_used_bytes BIGINT NOT NULL DEFAULT 0, + storage_total_bytes BIGINT NOT NULL DEFAULT 0, + disk_read_bytes BIGINT NOT NULL DEFAULT 0, + disk_write_bytes BIGINT NOT NULL DEFAULT 0, + network_rx_bytes BIGINT NOT NULL DEFAULT 0, + network_tx_bytes BIGINT NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_metric_snapshots_node_id ON of_node_metric_snapshots (node_id); +CREATE INDEX idx_of_node_metric_snapshots_captured_at ON of_node_metric_snapshots (captured_at); + +CREATE TABLE of_node_request_reports ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + window_started_at TIMESTAMPTZ NOT NULL, + window_ended_at TIMESTAMPTZ NOT NULL, + request_count BIGINT NOT NULL DEFAULT 0, + error_count BIGINT NOT NULL DEFAULT 0, + unique_visitor_count BIGINT NOT NULL DEFAULT 0, + status_codes_json TEXT, + top_domains_json TEXT, + source_countries_json TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_request_reports_node_id ON of_node_request_reports (node_id); +CREATE INDEX idx_of_node_request_reports_window_started_at ON of_node_request_reports (window_started_at); +CREATE INDEX idx_of_node_request_reports_window_ended_at ON of_node_request_reports (window_ended_at); + +CREATE TABLE of_node_health_events ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + event_type VARCHAR(64) NOT NULL, + severity VARCHAR(16) NOT NULL, + status VARCHAR(16) NOT NULL, + message TEXT, + first_triggered_at TIMESTAMPTZ NOT NULL, + last_triggered_at TIMESTAMPTZ NOT NULL, + reported_at TIMESTAMPTZ NOT NULL, + resolved_at TIMESTAMPTZ, + metadata_json TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_health_events_node_id ON of_node_health_events (node_id); +CREATE INDEX idx_of_node_health_events_event_type ON of_node_health_events (event_type); +CREATE INDEX idx_of_node_health_events_status ON of_node_health_events (status); +CREATE INDEX idx_of_node_health_events_first_triggered_at ON of_node_health_events (first_triggered_at); +CREATE INDEX idx_of_node_health_events_last_triggered_at ON of_node_health_events (last_triggered_at); +CREATE INDEX idx_of_node_health_events_reported_at ON of_node_health_events (reported_at); +CREATE INDEX idx_of_node_health_events_resolved_at ON of_node_health_events (resolved_at); + +CREATE TABLE of_node_obs_openresty ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + captured_at TIMESTAMPTZ NOT NULL, + openresty_rx_bytes BIGINT NOT NULL DEFAULT 0, + openresty_tx_bytes BIGINT NOT NULL DEFAULT 0, + openresty_connections BIGINT NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_obs_openresty_node_id ON of_node_obs_openresty (node_id); +CREATE INDEX idx_of_node_obs_openresty_captured_at ON of_node_obs_openresty (captured_at); + +CREATE TABLE of_node_obs_frps ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + captured_at TIMESTAMPTZ NOT NULL, + frps_connections INTEGER NOT NULL DEFAULT 0, + frps_proxy_count INTEGER NOT NULL DEFAULT 0, + frps_client_count INTEGER NOT NULL DEFAULT 0, + frps_proxies TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_obs_frps_node_id ON of_node_obs_frps (node_id); +CREATE INDEX idx_of_node_obs_frps_captured_at ON of_node_obs_frps (captured_at); + +CREATE TABLE of_node_access_logs ( + id BIGSERIAL PRIMARY KEY, + node_id VARCHAR(64) NOT NULL, + logged_at TIMESTAMPTZ NOT NULL, + remote_addr VARCHAR(128) NOT NULL DEFAULT '', + region VARCHAR(128) NOT NULL DEFAULT '', + host VARCHAR(255) NOT NULL DEFAULT '', + path VARCHAR(2048) NOT NULL DEFAULT '', + status_code INTEGER NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX idx_of_node_access_logs_node_id ON of_node_access_logs (node_id); +CREATE INDEX idx_of_node_access_logs_logged_at ON of_node_access_logs (logged_at); +CREATE INDEX idx_of_node_access_logs_remote_addr ON of_node_access_logs (remote_addr); +CREATE INDEX idx_of_node_access_logs_host ON of_node_access_logs (host); +CREATE INDEX idx_of_node_access_logs_status_code ON of_node_access_logs (status_code); + +-- +goose Down +DROP TABLE IF EXISTS of_node_access_logs; +DROP TABLE IF EXISTS of_node_obs_frps; +DROP TABLE IF EXISTS of_node_obs_openresty; +DROP TABLE IF EXISTS of_node_health_events; +DROP TABLE IF EXISTS of_node_request_reports; +DROP TABLE IF EXISTS of_node_metric_snapshots; +DROP TABLE IF EXISTS of_node_system_profiles; \ No newline at end of file diff --git a/Wavelet/internal/db/migrator/goose/sqlite/202606190010_create_of_observability_tables.sql b/Wavelet/internal/db/migrator/goose/sqlite/202606190010_create_of_observability_tables.sql new file mode 100644 index 00000000..bedc07de --- /dev/null +++ b/Wavelet/internal/db/migrator/goose/sqlite/202606190010_create_of_observability_tables.sql @@ -0,0 +1,136 @@ +-- +goose Up +CREATE TABLE IF NOT EXISTS of_node_system_profiles ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + hostname TEXT NOT NULL DEFAULT '', + os_name TEXT NOT NULL DEFAULT '', + os_version TEXT NOT NULL DEFAULT '', + kernel_version TEXT NOT NULL DEFAULT '', + architecture TEXT NOT NULL DEFAULT '', + cpu_model TEXT NOT NULL DEFAULT '', + cpu_cores INTEGER NOT NULL DEFAULT 0, + total_memory_bytes INTEGER NOT NULL DEFAULT 0, + total_disk_bytes INTEGER NOT NULL DEFAULT 0, + uptime_seconds INTEGER NOT NULL DEFAULT 0, + reported_at DATETIME NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_of_node_system_profiles_node_id ON of_node_system_profiles (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_system_profiles_reported_at ON of_node_system_profiles (reported_at); + +CREATE TABLE IF NOT EXISTS of_node_metric_snapshots ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + captured_at DATETIME NOT NULL, + cpu_usage_percent REAL NOT NULL DEFAULT 0, + memory_used_bytes INTEGER NOT NULL DEFAULT 0, + memory_total_bytes INTEGER NOT NULL DEFAULT 0, + storage_used_bytes INTEGER NOT NULL DEFAULT 0, + storage_total_bytes INTEGER NOT NULL DEFAULT 0, + disk_read_bytes INTEGER NOT NULL DEFAULT 0, + disk_write_bytes INTEGER NOT NULL DEFAULT 0, + network_rx_bytes INTEGER NOT NULL DEFAULT 0, + network_tx_bytes INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_metric_snapshots_node_id ON of_node_metric_snapshots (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_metric_snapshots_captured_at ON of_node_metric_snapshots (captured_at); + +CREATE TABLE IF NOT EXISTS of_node_request_reports ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + window_started_at DATETIME NOT NULL, + window_ended_at DATETIME NOT NULL, + request_count INTEGER NOT NULL DEFAULT 0, + error_count INTEGER NOT NULL DEFAULT 0, + unique_visitor_count INTEGER NOT NULL DEFAULT 0, + status_codes_json TEXT, + top_domains_json TEXT, + source_countries_json TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_request_reports_node_id ON of_node_request_reports (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_request_reports_window_started_at ON of_node_request_reports (window_started_at); +CREATE INDEX IF NOT EXISTS idx_of_node_request_reports_window_ended_at ON of_node_request_reports (window_ended_at); + +CREATE TABLE IF NOT EXISTS of_node_health_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + event_type TEXT NOT NULL, + severity TEXT NOT NULL, + status TEXT NOT NULL, + message TEXT, + first_triggered_at DATETIME NOT NULL, + last_triggered_at DATETIME NOT NULL, + reported_at DATETIME NOT NULL, + resolved_at DATETIME, + metadata_json TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_health_events_node_id ON of_node_health_events (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_health_events_event_type ON of_node_health_events (event_type); +CREATE INDEX IF NOT EXISTS idx_of_node_health_events_status ON of_node_health_events (status); +CREATE INDEX IF NOT EXISTS idx_of_node_health_events_first_triggered_at ON of_node_health_events (first_triggered_at); +CREATE INDEX IF NOT EXISTS idx_of_node_health_events_last_triggered_at ON of_node_health_events (last_triggered_at); +CREATE INDEX IF NOT EXISTS idx_of_node_health_events_reported_at ON of_node_health_events (reported_at); +CREATE INDEX IF NOT EXISTS idx_of_node_health_events_resolved_at ON of_node_health_events (resolved_at); + +CREATE TABLE IF NOT EXISTS of_node_obs_openresty ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + captured_at DATETIME NOT NULL, + openresty_rx_bytes INTEGER NOT NULL DEFAULT 0, + openresty_tx_bytes INTEGER NOT NULL DEFAULT 0, + openresty_connections INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_obs_openresty_node_id ON of_node_obs_openresty (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_obs_openresty_captured_at ON of_node_obs_openresty (captured_at); + +CREATE TABLE IF NOT EXISTS of_node_obs_frps ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + captured_at DATETIME NOT NULL, + frps_connections INTEGER NOT NULL DEFAULT 0, + frps_proxy_count INTEGER NOT NULL DEFAULT 0, + frps_client_count INTEGER NOT NULL DEFAULT 0, + frps_proxies TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_obs_frps_node_id ON of_node_obs_frps (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_obs_frps_captured_at ON of_node_obs_frps (captured_at); + +CREATE TABLE IF NOT EXISTS of_node_access_logs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id TEXT NOT NULL, + logged_at DATETIME NOT NULL, + remote_addr TEXT NOT NULL DEFAULT '', + region TEXT NOT NULL DEFAULT '', + host TEXT NOT NULL DEFAULT '', + path TEXT NOT NULL DEFAULT '', + status_code INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_node_id ON of_node_access_logs (node_id); +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_logged_at ON of_node_access_logs (logged_at); +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_remote_addr ON of_node_access_logs (remote_addr); +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_host ON of_node_access_logs (host); +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_status_code ON of_node_access_logs (status_code); + +-- +goose Down +DROP TABLE IF EXISTS of_node_access_logs; +DROP TABLE IF EXISTS of_node_obs_frps; +DROP TABLE IF EXISTS of_node_obs_openresty; +DROP TABLE IF EXISTS of_node_health_events; +DROP TABLE IF EXISTS of_node_request_reports; +DROP TABLE IF EXISTS of_node_metric_snapshots; +DROP TABLE IF EXISTS of_node_system_profiles; \ No newline at end of file diff --git a/Wavelet/internal/db/postgres.go b/Wavelet/internal/db/postgres.go index 7e01fba8..7aa72320 100644 --- a/Wavelet/internal/db/postgres.go +++ b/Wavelet/internal/db/postgres.go @@ -39,7 +39,7 @@ func init() { func initSQLite() { sqlitePath := config.Config.Database.SQLitePath if sqlitePath == "" { - sqlitePath = "./data/wavelet.db" + sqlitePath = "./data/openflare.db" } var err error diff --git a/Wavelet/support-files/sql/create_clickhouse_risk.sql b/Wavelet/support-files/sql/create_clickhouse_risk.sql index 412285b8..0d6a7664 100644 --- a/Wavelet/support-files/sql/create_clickhouse_risk.sql +++ b/Wavelet/support-files/sql/create_clickhouse_risk.sql @@ -1,6 +1,6 @@ -CREATE DATABASE IF NOT EXISTS wavelet; +CREATE DATABASE IF NOT EXISTS openflare; -USE wavelet; +USE openflare; CREATE TABLE IF NOT EXISTS w_user_access_logs ( diff --git a/docker-compose.yaml b/docker-compose.yaml index 63e3bff2..34ff0fca 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -16,7 +16,7 @@ services: environment: OPENFLARE_SERVER_URL: "http://host.docker.internal:3000" - OPENFLARE_AGENT_TOKEN: "07800f31d3f181e65d18dca1407d821c" + OPENFLARE_AGENT_TOKEN: "111675c3b22c8cf39b43639b9305fa0a" LOG_LEVEL: "debug" extra_hosts: diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 604ad48b..abecc724 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -23,6 +23,9 @@ sidebar: false - 重叠职能复用 Wavelet 内置用户/OAuth/Cap/认证源能力;新增 `integration` 包覆盖认证、核心链路、安全、Agent 协议集成测试。 - 修复 Wavelet 引入 `openflare` 模块后 `gomodule/redigo v2.0.0+incompatible` 导致 `redistore` 编译失败的问题。 - Wavelet 后端默认监听端口由 `:8000` 调整为 `:3000`,与旧 OpenFlare Server 保持一致。 +- 新增 OpenFlare 可观测性表 goose 迁移(`of_node_system_profiles`、`of_node_metric_snapshots`、`of_node_request_reports`、`of_node_health_events`、`of_node_obs_openresty`、`of_node_obs_frps`、`of_node_access_logs`)。 +- Agent heartbeat 恢复可观测性数据持久化(系统画像、指标快照、流量报表、健康事件等)。 +- Wavelet 默认数据库名由 `wavelet` 调整为 `openflare`(PostgreSQL、ClickHouse、SQLite 后备路径同步更新)。 ## [v2.3.4] - 2026-06-17