From d9d02e749ad829ac3ed3ec1b2d9a4f41023956b3 Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 14 Mar 2026 10:59:42 +0800 Subject: [PATCH] feat: add phase-one node observability ingestion --- atsf_agent/internal/agent/runner.go | 9 +- atsf_agent/internal/agent/runner_test.go | 44 ++ .../internal/observability/collector.go | 391 ++++++++++++++++++ atsf_agent/internal/protocol/agent_api.go | 71 +++- atsf_agent/internal/state/state.go | 16 +- atsf_server/model/main.go | 18 +- atsf_server/model/node_health_event.go | 37 ++ atsf_server/model/node_metric_snapshot.go | 39 ++ atsf_server/model/node_request_report.go | 34 ++ atsf_server/model/node_system_profile.go | 56 +++ atsf_server/service/agent.go | 23 +- atsf_server/service/node_update_test.go | 168 ++++++++ atsf_server/service/observability.go | 272 ++++++++++++ 13 files changed, 1152 insertions(+), 26 deletions(-) create mode 100644 atsf_agent/internal/observability/collector.go create mode 100644 atsf_server/model/node_health_event.go create mode 100644 atsf_server/model/node_metric_snapshot.go create mode 100644 atsf_server/model/node_request_report.go create mode 100644 atsf_server/model/node_system_profile.go create mode 100644 atsf_server/service/observability.go diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index 2ee6f881..65a24e89 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -8,6 +8,7 @@ import ( "time" "atsflare-agent/internal/config" + "atsflare-agent/internal/observability" "atsflare-agent/internal/protocol" "atsflare-agent/internal/state" ) @@ -232,7 +233,7 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error { r.recordSyncError(err) slog.Error("agent post-register startup sync failed", "error", err) } else { - slog.Debug("agent post-register startup sync completed") + slog.Debug("agent post-register startup sync completed") } r.tryRestartOpenresty(ctx) r.tryAutoUpdate(ctx) @@ -310,6 +311,9 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { if openrestyStatus == "" { openrestyStatus = protocol.OpenrestyStatusUnknown } + profile := observability.BuildProfile(r.Config, r.StateStore) + metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore) + healthEvents := observability.BuildHealthEvents(snapshot) return protocol.NodePayload{ NodeID: nodeID, Name: r.Config.NodeName, @@ -320,5 +324,8 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { LastError: snapshot.LastError, OpenrestyStatus: openrestyStatus, OpenrestyMessage: snapshot.OpenrestyMessage, + Profile: profile, + Snapshot: metricSnapshot, + HealthEvents: healthEvents, } } diff --git a/atsf_agent/internal/agent/runner_test.go b/atsf_agent/internal/agent/runner_test.go index c4e35c67..c572369c 100644 --- a/atsf_agent/internal/agent/runner_test.go +++ b/atsf_agent/internal/agent/runner_test.go @@ -281,6 +281,50 @@ func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) { } } +func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) { + stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json")) + if err := stateStore.Save(&state.Snapshot{ + NodeID: "node-observe", + CurrentVersion: "20260314-001", + LastError: "sync failed", + OpenrestyStatus: protocol.OpenrestyStatusUnhealthy, + OpenrestyMessage: "reload failed", + }); err != nil { + t.Fatalf("failed to seed state: %v", err) + } + + runner := &Runner{ + Config: &config.Config{ + NodeName: "edge-observe-1", + NodeIP: "10.0.0.51", + AgentVersion: config.AgentVersion, + NginxVersion: "1.27.1.2", + DataDir: t.TempDir(), + HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), + }, + StateStore: stateStore, + } + + firstPayload := runner.nodePayload("node-observe") + if firstPayload.Profile == nil { + t.Fatal("expected first heartbeat payload to include system profile") + } + if firstPayload.Snapshot == nil { + t.Fatal("expected first heartbeat payload to include metric snapshot") + } + if len(firstPayload.HealthEvents) != 2 { + t.Fatalf("expected health events for openresty and sync error, got %+v", firstPayload.HealthEvents) + } + + secondPayload := runner.nodePayload("node-observe") + if secondPayload.Profile != nil { + t.Fatal("expected unchanged profile to be omitted on subsequent heartbeat") + } + if secondPayload.Snapshot == nil { + t.Fatal("expected metric snapshot to continue reporting on subsequent heartbeat") + } +} + func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() diff --git a/atsf_agent/internal/observability/collector.go b/atsf_agent/internal/observability/collector.go new file mode 100644 index 00000000..58273a39 --- /dev/null +++ b/atsf_agent/internal/observability/collector.go @@ -0,0 +1,391 @@ +package observability + +import ( + "atsflare-agent/internal/config" + "atsflare-agent/internal/protocol" + "atsflare-agent/internal/state" + "bufio" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "os" + "path/filepath" + "runtime" + "strconv" + "strings" + "syscall" + "time" +) + +func BuildProfile(cfg *config.Config, stateStore *state.Store) *protocol.NodeSystemProfile { + profile := collectProfile(cfg) + if profile == nil { + return nil + } + fingerprint := fingerprintProfile(profile) + if stateStore == nil { + return profile + } + snapshot, err := stateStore.Load() + if err != nil { + return profile + } + if snapshot.LastProfileFingerprint == fingerprint { + return nil + } + snapshot.LastProfileFingerprint = fingerprint + if err = stateStore.Save(snapshot); err != nil { + return profile + } + return profile +} + +func BuildSnapshot(cfg *config.Config, stateStore *state.Store) *protocol.NodeMetricSnapshot { + now := time.Now().UTC() + metric := &protocol.NodeMetricSnapshot{ + CapturedAtUnix: now.Unix(), + } + + memTotal, memUsed := readMemInfo() + metric.MemoryTotalBytes = memTotal + metric.MemoryUsedBytes = memUsed + + storageTotal, storageUsed := statFilesystem(cfg.DataDir) + metric.StorageTotalBytes = storageTotal + metric.StorageUsedBytes = storageUsed + + metric.NetworkRxBytes, metric.NetworkTxBytes = readLinuxNetworkTotals() + metric.DiskReadBytes, metric.DiskWriteBytes = readLinuxDiskTotals() + + if stateStore == nil { + return metric + } + + totalCPU, idleCPU := readLinuxCPUStat() + snapshot, err := stateStore.Load() + if err != nil { + return metric + } + if snapshot.LastCPUStatTotal > 0 && totalCPU > snapshot.LastCPUStatTotal && idleCPU >= snapshot.LastCPUStatIdle { + deltaTotal := totalCPU - snapshot.LastCPUStatTotal + deltaIdle := idleCPU - snapshot.LastCPUStatIdle + if deltaTotal > 0 && deltaIdle <= deltaTotal { + metric.CPUUsagePercent = (float64(deltaTotal-deltaIdle) / float64(deltaTotal)) * 100 + } + } + snapshot.LastCPUStatTotal = totalCPU + snapshot.LastCPUStatIdle = idleCPU + snapshot.LastMetricAtUnix = now.Unix() + _ = stateStore.Save(snapshot) + + return metric +} + +func BuildHealthEvents(snapshot *state.Snapshot) []protocol.NodeHealthEvent { + if snapshot == nil { + return []protocol.NodeHealthEvent{} + } + events := make([]protocol.NodeHealthEvent, 0, 2) + nowUnix := time.Now().UTC().Unix() + if strings.TrimSpace(snapshot.OpenrestyStatus) == protocol.OpenrestyStatusUnhealthy { + events = append(events, protocol.NodeHealthEvent{ + EventType: "openresty_unhealthy", + Severity: "critical", + Message: strings.TrimSpace(snapshot.OpenrestyMessage), + TriggeredAtUnix: nowUnix, + }) + } + if strings.TrimSpace(snapshot.LastError) != "" { + events = append(events, protocol.NodeHealthEvent{ + EventType: "sync_error", + Severity: "warning", + Message: strings.TrimSpace(snapshot.LastError), + TriggeredAtUnix: nowUnix, + }) + } + return events +} + +func collectProfile(cfg *config.Config) *protocol.NodeSystemProfile { + hostname, _ := os.Hostname() + osName, osVersion := readLinuxOSRelease() + kernelVersion := readFirstLine("/proc/sys/kernel/osrelease") + cpuModel := readLinuxCPUModel() + totalMemory, _ := readMemInfo() + totalDisk, _ := statFilesystem(cfg.DataDir) + uptimeSeconds := readLinuxUptimeSeconds() + + return &protocol.NodeSystemProfile{ + Hostname: strings.TrimSpace(hostname), + OSName: osName, + OSVersion: osVersion, + KernelVersion: kernelVersion, + Architecture: runtime.GOARCH, + CPUModel: cpuModel, + CPUCores: runtime.NumCPU(), + TotalMemoryBytes: totalMemory, + TotalDiskBytes: totalDisk, + UptimeSeconds: uptimeSeconds, + ReportedAtUnix: time.Now().UTC().Unix(), + } +} + +func fingerprintProfile(profile *protocol.NodeSystemProfile) string { + raw, err := json.Marshal(profile) + if err != nil { + return "" + } + sum := sha256.Sum256(raw) + return hex.EncodeToString(sum[:]) +} + +func readLinuxOSRelease() (string, string) { + file, err := os.Open("/etc/os-release") + if err != nil { + return runtime.GOOS, "" + } + defer file.Close() + + values := make(map[string]string) + scanner := bufio.NewScanner(file) + for scanner.Scan() { + line := strings.TrimSpace(scanner.Text()) + if line == "" || strings.HasPrefix(line, "#") { + continue + } + key, value, ok := strings.Cut(line, "=") + if !ok { + continue + } + values[key] = strings.Trim(value, `"`) + } + if pretty := strings.TrimSpace(values["PRETTY_NAME"]); pretty != "" { + return pretty, strings.TrimSpace(values["VERSION_ID"]) + } + name := strings.TrimSpace(values["NAME"]) + if name == "" { + name = runtime.GOOS + } + return name, strings.TrimSpace(values["VERSION_ID"]) +} + +func readLinuxCPUModel() string { + file, err := os.Open("/proc/cpuinfo") + if err != nil { + return "" + } + defer file.Close() + + scanner := bufio.NewScanner(file) + for scanner.Scan() { + line := scanner.Text() + if strings.HasPrefix(strings.ToLower(line), "model name") { + _, value, ok := strings.Cut(line, ":") + if ok { + return strings.TrimSpace(value) + } + } + } + return "" +} + +func readMemInfo() (int64, int64) { + file, err := os.Open("/proc/meminfo") + if err != nil { + return 0, 0 + } + defer file.Close() + + var memTotalKB int64 + var memAvailableKB int64 + scanner := bufio.NewScanner(file) + for scanner.Scan() { + line := scanner.Text() + switch { + case strings.HasPrefix(line, "MemTotal:"): + memTotalKB = parseMemInfoValue(line) + case strings.HasPrefix(line, "MemAvailable:"): + memAvailableKB = parseMemInfoValue(line) + } + } + + total := memTotalKB * 1024 + if total == 0 { + return 0, 0 + } + used := total - (memAvailableKB * 1024) + if used < 0 { + used = 0 + } + return total, used +} + +func parseMemInfoValue(line string) int64 { + fields := strings.Fields(line) + if len(fields) < 2 { + return 0 + } + value, err := strconv.ParseInt(fields[1], 10, 64) + if err != nil { + return 0 + } + return value +} + +func readLinuxUptimeSeconds() int64 { + content, err := os.ReadFile("/proc/uptime") + if err != nil { + return 0 + } + fields := strings.Fields(string(content)) + if len(fields) == 0 { + return 0 + } + value, err := strconv.ParseFloat(fields[0], 64) + if err != nil { + return 0 + } + return int64(value) +} + +func readLinuxCPUStat() (uint64, uint64) { + content, err := os.ReadFile("/proc/stat") + if err != nil { + return 0, 0 + } + lines := strings.Split(string(content), "\n") + for _, line := range lines { + if !strings.HasPrefix(line, "cpu ") { + continue + } + fields := strings.Fields(line) + if len(fields) < 5 { + return 0, 0 + } + var total uint64 + for i := 1; i < len(fields); i++ { + value, err := strconv.ParseUint(fields[i], 10, 64) + if err != nil { + return 0, 0 + } + total += value + if i == 4 { + // idle + } + } + idle, err := strconv.ParseUint(fields[4], 10, 64) + if err != nil { + return 0, 0 + } + return total, idle + } + return 0, 0 +} + +func readLinuxNetworkTotals() (int64, int64) { + file, err := os.Open("/proc/net/dev") + if err != nil { + return 0, 0 + } + defer file.Close() + + var rx int64 + var tx int64 + scanner := bufio.NewScanner(file) + for scanner.Scan() { + line := strings.TrimSpace(scanner.Text()) + if !strings.Contains(line, ":") { + continue + } + name, data, ok := strings.Cut(line, ":") + if !ok { + continue + } + if strings.TrimSpace(name) == "lo" { + continue + } + fields := strings.Fields(data) + if len(fields) < 16 { + continue + } + rxValue, err := strconv.ParseInt(fields[0], 10, 64) + if err == nil { + rx += rxValue + } + txValue, err := strconv.ParseInt(fields[8], 10, 64) + if err == nil { + tx += txValue + } + } + return rx, tx +} + +func readLinuxDiskTotals() (int64, int64) { + file, err := os.Open("/proc/diskstats") + if err != nil { + return 0, 0 + } + defer file.Close() + + var readBytes int64 + var writeBytes int64 + scanner := bufio.NewScanner(file) + for scanner.Scan() { + fields := strings.Fields(scanner.Text()) + if len(fields) < 14 { + continue + } + device := fields[2] + if shouldSkipDiskDevice(device) { + continue + } + readSectors, err := strconv.ParseInt(fields[5], 10, 64) + if err == nil { + readBytes += readSectors * 512 + } + writeSectors, err := strconv.ParseInt(fields[9], 10, 64) + if err == nil { + writeBytes += writeSectors * 512 + } + } + return readBytes, writeBytes +} + +func shouldSkipDiskDevice(device string) bool { + switch { + case device == "": + return true + case strings.HasPrefix(device, "loop"), + strings.HasPrefix(device, "ram"), + strings.HasPrefix(device, "dm-"): + return true + default: + return false + } +} + +func statFilesystem(path string) (int64, int64) { + if strings.TrimSpace(path) == "" { + path = string(os.PathSeparator) + } + absPath := filepath.Clean(path) + var stat syscall.Statfs_t + if err := syscall.Statfs(absPath, &stat); err != nil { + return 0, 0 + } + total := int64(stat.Blocks) * int64(stat.Bsize) + free := int64(stat.Bavail) * int64(stat.Bsize) + used := total - free + if used < 0 { + used = 0 + } + return total, used +} + +func readFirstLine(path string) string { + content, err := os.ReadFile(path) + if err != nil { + return "" + } + return strings.TrimSpace(string(content)) +} diff --git a/atsf_agent/internal/protocol/agent_api.go b/atsf_agent/internal/protocol/agent_api.go index 931df306..996a2727 100644 --- a/atsf_agent/internal/protocol/agent_api.go +++ b/atsf_agent/internal/protocol/agent_api.go @@ -36,15 +36,68 @@ const ( ) 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"` - CurrentVersion string `json:"current_version"` - LastError string `json:"last_error"` - OpenrestyStatus string `json:"openresty_status"` - OpenrestyMessage string `json:"openresty_message"` + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + AgentVersion string `json:"agent_version"` + NginxVersion 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"` + TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"` + HealthEvents []NodeHealthEvent `json:"health_events"` +} + +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"` +} + +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"` +} + +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"` +} + +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,omitempty"` } type RegisterNodeResponse struct { diff --git a/atsf_agent/internal/state/state.go b/atsf_agent/internal/state/state.go index 8f7b101e..7fc60254 100644 --- a/atsf_agent/internal/state/state.go +++ b/atsf_agent/internal/state/state.go @@ -10,12 +10,16 @@ import ( ) type Snapshot struct { - NodeID string `json:"node_id"` - CurrentVersion string `json:"current_version"` - CurrentChecksum string `json:"current_checksum"` - LastError string `json:"last_error"` - OpenrestyStatus string `json:"openresty_status"` - OpenrestyMessage string `json:"openresty_message"` + NodeID string `json:"node_id"` + CurrentVersion string `json:"current_version"` + CurrentChecksum string `json:"current_checksum"` + LastError string `json:"last_error"` + OpenrestyStatus string `json:"openresty_status"` + OpenrestyMessage string `json:"openresty_message"` + LastProfileFingerprint string `json:"last_profile_fingerprint"` + LastCPUStatTotal uint64 `json:"last_cpu_stat_total"` + LastCPUStatIdle uint64 `json:"last_cpu_stat_idle"` + LastMetricAtUnix int64 `json:"last_metric_at_unix"` } type Store struct { diff --git a/atsf_server/model/main.go b/atsf_server/model/main.go index 35610830..93eb300c 100644 --- a/atsf_server/model/main.go +++ b/atsf_server/model/main.go @@ -3,9 +3,9 @@ package model import ( "atsflare/common" "github.com/glebarez/sqlite" - "log/slog" "gorm.io/driver/mysql" "gorm.io/gorm" + "log/slog" "os" ) @@ -90,10 +90,26 @@ func InitDB() (err error) { if err != nil { return err } + err = db.AutoMigrate(&NodeSystemProfile{}) + if err != nil { + return err + } err = db.AutoMigrate(&ApplyLog{}) if err != nil { return err } + err = db.AutoMigrate(&NodeMetricSnapshot{}) + if err != nil { + return err + } + err = db.AutoMigrate(&NodeRequestReport{}) + if err != nil { + return err + } + err = db.AutoMigrate(&NodeHealthEvent{}) + if err != nil { + return err + } err = db.AutoMigrate(&TLSCertificate{}) if err != nil { return err diff --git a/atsf_server/model/node_health_event.go b/atsf_server/model/node_health_event.go new file mode 100644 index 00000000..fade3d32 --- /dev/null +++ b/atsf_server/model/node_health_event.go @@ -0,0 +1,37 @@ +package model + +import "time" + +type NodeHealthEvent struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"index;size:64;not null"` + EventType string `json:"event_type" gorm:"index;size:64;not null"` + Severity string `json:"severity" gorm:"size:16;not null"` + Status string `json:"status" gorm:"index;size:16;not null"` + Message string `json:"message" gorm:"size:2048"` + FirstTriggeredAt time.Time `json:"first_triggered_at" gorm:"index"` + 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"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +func GetActiveNodeHealthEvent(nodeID string, eventType string) (*NodeHealthEvent, error) { + event := &NodeHealthEvent{} + err := DB.Where("node_id = ? AND event_type = ? AND status = ?", nodeID, eventType, "active").First(event).Error + return event, err +} + +func ListNodeHealthEvents(nodeID string, activeOnly bool, limit int) (events []*NodeHealthEvent, err error) { + query := DB.Where("node_id = ?", nodeID).Order("last_triggered_at desc") + if activeOnly { + query = query.Where("status = ?", "active") + } + if limit > 0 { + query = query.Limit(limit) + } + err = query.Find(&events).Error + return events, err +} diff --git a/atsf_server/model/node_metric_snapshot.go b/atsf_server/model/node_metric_snapshot.go new file mode 100644 index 00000000..766ddefe --- /dev/null +++ b/atsf_server/model/node_metric_snapshot.go @@ -0,0 +1,39 @@ +package model + +import "time" + +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"` + RawJSON string `json:"raw_json" gorm:"type:text"` + CreatedAt time.Time `json:"created_at"` +} + +func (snapshot *NodeMetricSnapshot) Insert() error { + return DB.Create(snapshot).Error +} + +func ListNodeMetricSnapshots(nodeID string, since time.Time, limit int) (snapshots []*NodeMetricSnapshot, err error) { + query := DB.Where("node_id = ?", nodeID).Order("captured_at desc") + if !since.IsZero() { + query = query.Where("captured_at >= ?", since) + } + if limit > 0 { + query = query.Limit(limit) + } + err = query.Find(&snapshots).Error + return snapshots, err +} diff --git a/atsf_server/model/node_request_report.go b/atsf_server/model/node_request_report.go new file mode 100644 index 00000000..f475bc6a --- /dev/null +++ b/atsf_server/model/node_request_report.go @@ -0,0 +1,34 @@ +package model + +import "time" + +type NodeRequestReport struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"index;size:64;not null"` + WindowStartedAt time.Time `json:"window_started_at" gorm:"index"` + WindowEndedAt time.Time `json:"window_ended_at" gorm:"index"` + RequestCount int64 `json:"request_count"` + ErrorCount int64 `json:"error_count"` + UniqueVisitorCount int64 `json:"unique_visitor_count"` + 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"` +} + +func (report *NodeRequestReport) Insert() error { + return DB.Create(report).Error +} + +func ListNodeRequestReports(nodeID string, since time.Time, limit int) (reports []*NodeRequestReport, err error) { + query := DB.Where("node_id = ?", nodeID).Order("window_ended_at desc") + if !since.IsZero() { + query = query.Where("window_ended_at >= ?", since) + } + if limit > 0 { + query = query.Limit(limit) + } + err = query.Find(&reports).Error + return reports, err +} diff --git a/atsf_server/model/node_system_profile.go b/atsf_server/model/node_system_profile.go new file mode 100644 index 00000000..01866aeb --- /dev/null +++ b/atsf_server/model/node_system_profile.go @@ -0,0 +1,56 @@ +package model + +import ( + "time" + + "gorm.io/gorm/clause" +) + +type NodeSystemProfile struct { + ID uint `json:"id" gorm:"primaryKey"` + NodeID string `json:"node_id" gorm:"uniqueIndex;size:64;not null"` + Hostname string `json:"hostname" gorm:"size:255"` + OSName string `json:"os_name" gorm:"size:128"` + OSVersion string `json:"os_version" gorm:"size:128"` + KernelVersion string `json:"kernel_version" gorm:"size:128"` + Architecture string `json:"architecture" gorm:"size:64"` + CPUModel string `json:"cpu_model" gorm:"size:255"` + CPUCores int `json:"cpu_cores"` + TotalMemoryBytes int64 `json:"total_memory_bytes"` + 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"` +} + +func GetNodeSystemProfile(nodeID string) (*NodeSystemProfile, error) { + profile := &NodeSystemProfile{} + err := DB.Where("node_id = ?", nodeID).First(profile).Error + return profile, err +} + +func UpsertNodeSystemProfile(profile *NodeSystemProfile) error { + if profile == nil { + return nil + } + return DB.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", + "raw_json", + "updated_at", + }), + }).Create(profile).Error +} diff --git a/atsf_server/service/agent.go b/atsf_server/service/agent.go index cf29ee1b..769ce897 100644 --- a/atsf_server/service/agent.go +++ b/atsf_server/service/agent.go @@ -24,15 +24,19 @@ const ( ) 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"` - CurrentVersion string `json:"current_version"` - LastError string `json:"last_error"` - OpenrestyStatus string `json:"openresty_status"` - OpenrestyMessage string `json:"openresty_message"` + NodeID string `json:"node_id"` + Name string `json:"name"` + IP string `json:"ip"` + AgentVersion string `json:"agent_version"` + NginxVersion 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 *AgentNodeSystemProfile `json:"profile,omitempty"` + Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"` + TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"` + HealthEvents []AgentNodeHealthEvent `json:"health_events"` } type ApplyLogPayload struct { @@ -134,6 +138,7 @@ func HeartbeatNode(node *model.Node, payload AgentNodePayload) (*HeartbeatRespon return nil, err } } + persistHeartbeatObservability(node.NodeID, payload, node.LastSeenAt) activeConfig, err := GetActiveConfigMetaForAgent() if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { return nil, err diff --git a/atsf_server/service/node_update_test.go b/atsf_server/service/node_update_test.go index 6af5ea47..bdf330f9 100644 --- a/atsf_server/service/node_update_test.go +++ b/atsf_server/service/node_update_test.go @@ -325,3 +325,171 @@ func TestListNodeViewsDoesNotPersistComputedStatus(t *testing.T) { t.Fatalf("expected list query to avoid persisting computed status, got %s", storedNode.Status) } } + +func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) { + setupServiceTestDB(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, + } + 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, + Profile: &AgentNodeSystemProfile{ + Hostname: "observe-edge-1", + OSName: "Ubuntu", + OSVersion: "24.04", + KernelVersion: "6.8.0", + Architecture: "amd64", + CPUModel: "Intel Xeon", + CPUCores: 8, + TotalMemoryBytes: 16 * 1024 * 1024 * 1024, + TotalDiskBytes: 200 * 1024 * 1024 * 1024, + UptimeSeconds: 3600, + 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, + }, + TrafficReport: &AgentNodeTrafficReport{ + WindowStartedAtUnix: time.Now().Add(-time.Minute).Unix(), + WindowEndedAtUnix: time.Now().Unix(), + RequestCount: 1200, + ErrorCount: 12, + UniqueVisitorCount: 320, + StatusCodes: map[string]int64{"200": 1100, "502": 12}, + TopDomains: map[string]int64{"example.com": 900}, + SourceCountries: map[string]int64{"CN": 700, "US": 200}, + }, + HealthEvents: []AgentNodeHealthEvent{ + { + EventType: "openresty_unhealthy", + Severity: NodeHealthSeverityCritical, + Message: "reload failed", + TriggeredAtUnix: time.Now().Add(-2 * time.Minute).Unix(), + }, + }, + }) + if err != nil { + t.Fatalf("expected heartbeat to succeed: %v", err) + } + + profile, err := model.GetNodeSystemProfile(node.NodeID) + if err != nil { + t.Fatalf("expected node profile to persist: %v", err) + } + if profile.OSName != "Ubuntu" || profile.CPUCores != 8 { + t.Fatalf("unexpected system profile: %+v", profile) + } + + snapshots, err := model.ListNodeMetricSnapshots(node.NodeID, time.Time{}, 10) + if err != nil { + t.Fatalf("expected node snapshots query to succeed: %v", err) + } + if len(snapshots) != 1 || snapshots[0].OpenrestyConnections != 128 { + t.Fatalf("unexpected metric snapshots: %+v", snapshots) + } + + reports, err := model.ListNodeRequestReports(node.NodeID, time.Time{}, 10) + if err != nil { + t.Fatalf("expected node request reports query to succeed: %v", err) + } + if len(reports) != 1 || reports[0].RequestCount != 1200 { + t.Fatalf("unexpected request reports: %+v", reports) + } + + events, err := model.ListNodeHealthEvents(node.NodeID, true, 10) + if err != nil { + t.Fatalf("expected node health events query to succeed: %v", err) + } + if len(events) != 1 || events[0].EventType != "openresty_unhealthy" { + t.Fatalf("unexpected active health events: %+v", events) + } +} + +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, + } + 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, + HealthEvents: []AgentNodeHealthEvent{ + { + EventType: "sync_error", + Severity: NodeHealthSeverityWarning, + Message: "checksum mismatch", + TriggeredAtUnix: time.Now().Add(-time.Minute).Unix(), + }, + }, + }) + if err != nil { + t.Fatalf("expected first heartbeat to succeed: %v", err) + } + + _, err = HeartbeatNode(node, AgentNodePayload{ + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + AgentVersion: node.AgentVersion, + NginxVersion: node.NginxVersion, + HealthEvents: []AgentNodeHealthEvent{}, + }) + if err != nil { + t.Fatalf("expected second heartbeat to succeed: %v", err) + } + + activeEvents, err := model.ListNodeHealthEvents(node.NodeID, true, 10) + if err != nil { + t.Fatalf("expected active node health events query to succeed: %v", err) + } + if len(activeEvents) != 0 { + t.Fatalf("expected no active health events, got %+v", activeEvents) + } + + allEvents, err := model.ListNodeHealthEvents(node.NodeID, false, 10) + if err != nil { + t.Fatalf("expected all node health events query to succeed: %v", err) + } + if len(allEvents) != 1 || allEvents[0].Status != NodeHealthEventStatusResolved || allEvents[0].ResolvedAt == nil { + t.Fatalf("expected resolved health event record, got %+v", allEvents) + } +} diff --git a/atsf_server/service/observability.go b/atsf_server/service/observability.go new file mode 100644 index 00000000..770e9750 --- /dev/null +++ b/atsf_server/service/observability.go @@ -0,0 +1,272 @@ +package service + +import ( + "atsflare/model" + "encoding/json" + "errors" + "log/slog" + "strings" + "time" + + "gorm.io/gorm" +) + +const ( + NodeHealthEventStatusActive = "active" + NodeHealthEventStatusResolved = "resolved" + NodeHealthSeverityInfo = "info" + NodeHealthSeverityWarning = "warning" + NodeHealthSeverityCritical = "critical" +) + +type AgentNodeSystemProfile 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"` +} + +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"` +} + +type AgentNodeTrafficReport 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"` +} + +type AgentNodeHealthEvent 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(nodeID string, payload AgentNodePayload, reportedAt time.Time) { + if strings.TrimSpace(nodeID) == "" { + return + } + if payload.Profile == nil && payload.Snapshot == nil && payload.TrafficReport == nil && payload.HealthEvents == nil { + return + } + + if err := model.DB.Transaction(func(tx *gorm.DB) error { + if err := persistNodeSystemProfile(tx, nodeID, payload.Profile, reportedAt); err != nil { + return err + } + if err := persistNodeMetricSnapshot(tx, nodeID, payload.Snapshot, reportedAt); err != nil { + return err + } + if err := persistNodeTrafficReport(tx, nodeID, payload.TrafficReport, 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 { + slog.Error("persist heartbeat observability failed", "node_id", nodeID, "error", err) + } +} + +func persistNodeSystemProfile(tx *gorm.DB, nodeID string, profile *AgentNodeSystemProfile, reportedAt time.Time) error { + if profile == nil { + return nil + } + record := &model.NodeSystemProfile{ + 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), + RawJSON: marshalJSON(profile), + } + return tx.Model(&model.NodeSystemProfile{}).Where("node_id = ?", nodeID).Assign(record).FirstOrCreate(record).Error +} + +func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *AgentNodeMetricSnapshot, reportedAt time.Time) error { + if snapshot == nil { + 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, + RawJSON: marshalJSON(snapshot), + } + return tx.Create(record).Error +} + +func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *AgentNodeTrafficReport, 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.NodeRequestReport{ + 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), + RawJSON: marshalJSON(report), + } + return tx.Create(record).Error +} + +func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHealthEvent, reportedAt time.Time) error { + activeTypes := make(map[string]AgentNodeHealthEvent, 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.NodeHealthEvent + if err := tx.Where("node_id = ? AND status = ?", nodeID, NodeHealthEventStatusActive).Find(&activeEvents).Error; err != nil { + return err + } + + activeByType := make(map[string]*model.NodeHealthEvent, 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 = strings.TrimSpace(event.Message) + existing.LastTriggeredAt = triggeredAt + existing.ReportedAt = reportedAt + existing.RawJSON = marshalJSON(event) + existing.ResolvedAt = nil + if err := tx.Save(existing).Error; err != nil { + return err + } + continue + } + record := &model.NodeHealthEvent{ + NodeID: nodeID, + EventType: eventType, + Severity: event.Severity, + Status: NodeHealthEventStatusActive, + Message: strings.TrimSpace(event.Message), + FirstTriggeredAt: triggeredAt, + LastTriggeredAt: triggeredAt, + ReportedAt: reportedAt, + RawJSON: marshalJSON(event), + } + 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 = NodeHealthEventStatusResolved + existing.ReportedAt = reportedAt + existing.ResolvedAt = &resolvedAt + if err := tx.Save(existing).Error; err != nil { + return err + } + } + + return 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 NodeHealthSeverityCritical: + return NodeHealthSeverityCritical + case NodeHealthSeverityInfo: + return NodeHealthSeverityInfo + default: + return NodeHealthSeverityWarning + } +} + +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) +}