feat: add phase-one node observability ingestion

This commit is contained in:
ryan
2026-03-14 10:59:42 +08:00
parent 295340dc4b
commit d9d02e749a
13 changed files with 1152 additions and 26 deletions
+8 -1
View File
@@ -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,
}
}
+44
View File
@@ -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()
@@ -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))
}
+62 -9
View File
@@ -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 {
+10 -6
View File
@@ -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 {
+17 -1
View File
@@ -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
+37
View File
@@ -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
}
+39
View File
@@ -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
}
+34
View File
@@ -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
}
+56
View File
@@ -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
}
+14 -9
View File
@@ -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
+168
View File
@@ -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)
}
}
+272
View File
@@ -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)
}