diff --git a/README.md b/README.md index 61e5bdc3..c7b7121c 100644 --- a/README.md +++ b/README.md @@ -216,6 +216,8 @@ Docker 镜像工作流仅构建 `atsf_server`,并产出 `linux/amd64` 与 `lin | `data_dir` | Agent 托管数据目录 | | `support_dir` | Agent 存放受管附属文件的目录,当前包含证书与 Lua 观测脚本 | | `openresty_observability_port` | Agent 读取 OpenResty Lua 本地观测指标的 loopback 端口 | +| `observability_buffer_path` | Agent 本地观测补报缓冲文件路径 | +| `observability_replay_minutes` | Agent 恢复 heartbeat 后允许自动补传的最近观测窗口分钟数 | | `nginx_path` | 本机 Nginx 路径,设置后走本机模式 | | `nginx_container_name` | Docker 模式下的 Nginx 容器名 | diff --git a/atsf_agent/cmd/agent/main.go b/atsf_agent/cmd/agent/main.go index 34fd808f..99437518 100644 --- a/atsf_agent/cmd/agent/main.go +++ b/atsf_agent/cmd/agent/main.go @@ -55,6 +55,7 @@ func main() { client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration()) stateStore := state.NewStore(cfg.StatePath) + observabilityBuffer := state.NewObservabilityBufferStore(cfg.ObservabilityBufferPath) runtimeRouteConfigPath := cfg.RouteConfigPath if cfg.OpenrestyPath == "" { runtimeRouteConfigPath = nginx.DockerRouteConfigPath @@ -80,12 +81,13 @@ func main() { }), } runner := &agent.Runner{ - Config: cfg, - StateStore: stateStore, - HeartbeatService: heartbeat.New(client), - SyncService: syncservice.New(client, runtimeManager, stateStore), - Updater: updater.New(), - RuntimeManager: runtimeManager, + Config: cfg, + StateStore: stateStore, + ObservabilityBuffer: observabilityBuffer, + HeartbeatService: heartbeat.New(client), + SyncService: syncservice.New(client, runtimeManager, stateStore), + Updater: updater.New(), + RuntimeManager: runtimeManager, } ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index 755dcec4..e921a389 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -40,12 +40,13 @@ type UpdateOptions struct { } type Runner struct { - Config *config.Config - StateStore *state.Store - HeartbeatService HeartbeatService - SyncService SyncService - Updater Updater - RuntimeManager RuntimeManager + Config *config.Config + StateStore *state.Store + ObservabilityBuffer *state.ObservabilityBufferStore + HeartbeatService HeartbeatService + SyncService SyncService + Updater Updater + RuntimeManager RuntimeManager autoUpdate bool updateNow bool @@ -63,10 +64,12 @@ func (r *Runner) Run(ctx context.Context) error { slog.Info("agent runner started", "node_id", nodeID, "node", r.Config.NodeName, "ip", r.Config.NodeIP) if r.hasAgentToken() { r.refreshOpenrestyHealth(ctx) - heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) + payload, ackWindows := r.prepareHeartbeatPayload(nodeID) + heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, payload) if hbErr != nil { slog.Error("agent startup heartbeat failed", "error", hbErr) } else { + r.ackObservabilityWindows(ackWindows) if heartbeatResult == nil { heartbeatResult = &protocol.HeartbeatResult{} } @@ -101,10 +104,12 @@ func (r *Runner) Run(ctx context.Context) error { continue } r.refreshOpenrestyHealth(ctx) - heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) + payload, ackWindows := r.prepareHeartbeatPayload(nodeID) + heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, payload) if hbErr != nil { slog.Error("agent heartbeat failed", "error", hbErr) } else { + r.ackObservabilityWindows(ackWindows) if heartbeatResult == nil { heartbeatResult = &protocol.HeartbeatResult{} } @@ -220,11 +225,13 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error { *nodeID = response.NodeID slog.Info("agent discovery registration succeeded", "node_id", response.NodeID) r.refreshOpenrestyHealth(ctx) - heartbeatResult, heartbeatErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(*nodeID)) + payload, ackWindows := r.prepareHeartbeatPayload(*nodeID) + heartbeatResult, heartbeatErr := r.HeartbeatService.Heartbeat(ctx, payload) if heartbeatErr != nil { slog.Error("agent post-register heartbeat failed", "error", heartbeatErr) return nil } + r.ackObservabilityWindows(ackWindows) if heartbeatResult == nil { heartbeatResult = &protocol.HeartbeatResult{} } @@ -332,3 +339,60 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { HealthEvents: healthEvents, } } + +func (r *Runner) prepareHeartbeatPayload(nodeID string) (protocol.NodePayload, []int64) { + payload := r.nodePayload(nodeID) + if r.ObservabilityBuffer == nil || payload.Snapshot == nil { + return payload, nil + } + now := time.Now().UTC() + retainAfterUnix := now.Add(-time.Duration(r.Config.ObservabilityReplayMinutes) * time.Minute).Unix() + windowStartedAtUnix := state.ObservabilityWindowStartedAt(payload.Snapshot, payload.TrafficReport) + if windowStartedAtUnix <= 0 { + return payload, nil + } + + record := state.ObservabilityBufferRecord{ + WindowStartedAtUnix: windowStartedAtUnix, + Snapshot: payload.Snapshot, + TrafficReport: payload.TrafficReport, + QueuedAtUnix: now.Unix(), + } + if err := r.ObservabilityBuffer.Upsert(record, retainAfterUnix); err != nil { + slog.Error("upsert observability buffer failed", "error", err) + return payload, nil + } + + records, err := r.ObservabilityBuffer.Replayable(windowStartedAtUnix, retainAfterUnix) + if err != nil { + slog.Error("load replayable observability buffer failed", "error", err) + return payload, []int64{windowStartedAtUnix} + } + + ackWindows := make([]int64, 0, len(records)+1) + buffered := make([]protocol.BufferedObservabilityRecord, 0, len(records)) + for _, item := range records { + if item.WindowStartedAtUnix <= 0 { + continue + } + buffered = append(buffered, protocol.BufferedObservabilityRecord{ + WindowStartedAtUnix: item.WindowStartedAtUnix, + Snapshot: item.Snapshot, + TrafficReport: item.TrafficReport, + }) + ackWindows = append(ackWindows, item.WindowStartedAtUnix) + } + payload.BufferedObservability = buffered + ackWindows = append(ackWindows, windowStartedAtUnix) + return payload, ackWindows +} + +func (r *Runner) ackObservabilityWindows(windowStartedAtUnix []int64) { + if r.ObservabilityBuffer == nil || len(windowStartedAtUnix) == 0 { + return + } + retainAfterUnix := time.Now().UTC().Add(-time.Duration(r.Config.ObservabilityReplayMinutes) * time.Minute).Unix() + if err := r.ObservabilityBuffer.Ack(windowStartedAtUnix, retainAfterUnix); err != nil { + slog.Error("ack observability buffer failed", "error", err) + } +} diff --git a/atsf_agent/internal/agent/runner_test.go b/atsf_agent/internal/agent/runner_test.go index 5dec5516..2e2d4783 100644 --- a/atsf_agent/internal/agent/runner_test.go +++ b/atsf_agent/internal/agent/runner_test.go @@ -311,7 +311,7 @@ func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) { } if err := os.WriteFile( filepath.Join(filepath.Dir(runner.Config.RouteConfigPath), "atsflare_access.log"), - []byte("{\"ts\":\"2026-03-14T10:00:00Z\",\"host\":\"edge.example.com\",\"remote_addr\":\"10.0.0.8\",\"status\":200}\n"), + []byte("{\"ts\":\""+time.Now().UTC().Format(time.RFC3339)+"\",\"host\":\"edge.example.com\",\"remote_addr\":\"10.0.0.8\",\"status\":200}\n"), 0o644, ); err != nil { t.Fatalf("failed to prepare access log: %v", err) @@ -343,6 +343,81 @@ func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) { } } +func TestRunnerReplaysBufferedObservabilityAfterHeartbeatRecovery(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + tempDir := t.TempDir() + stateStore := state.NewStore(filepath.Join(tempDir, "state.json")) + bufferStore := state.NewObservabilityBufferStore(filepath.Join(tempDir, "observability-buffer.json")) + nowUnix := time.Now().UTC().Unix() + bufferWindow := nowUnix - (nowUnix % 60) - 60 + if err := bufferStore.Upsert(state.ObservabilityBufferRecord{ + WindowStartedAtUnix: bufferWindow, + Snapshot: &protocol.NodeMetricSnapshot{CapturedAtUnix: bufferWindow + 5, CPUUsagePercent: 30}, + TrafficReport: &protocol.NodeTrafficReport{WindowStartedAtUnix: bufferWindow, WindowEndedAtUnix: bufferWindow + 60, RequestCount: 8}, + QueuedAtUnix: bufferWindow + 60, + }, 0); err != nil { + t.Fatalf("failed to seed observability buffer: %v", err) + } + heartbeatService := &fakeHeartbeatService{ + heartbeatErrs: []error{errors.New("server offline"), nil}, + heartbeatResults: []*protocol.HeartbeatResult{{}, {}}, + onHeartbeat: func(callCount int) { + if callCount >= 2 { + cancel() + } + }, + } + runner := &Runner{ + Config: &config.Config{ + AgentToken: "agent-token", + NodeName: "edge-buffer-01", + NodeIP: "10.0.0.52", + AgentVersion: config.AgentVersion, + NginxVersion: "1.27.1.2", + DataDir: tempDir, + RouteConfigPath: filepath.Join(tempDir, "conf.d", "atsflare_routes.conf"), + HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond), + ObservabilityReplayMinutes: 15, + }, + StateStore: stateStore, + ObservabilityBuffer: bufferStore, + HeartbeatService: heartbeatService, + SyncService: &fakeSyncService{}, + } + if err := os.MkdirAll(filepath.Dir(runner.Config.RouteConfigPath), 0o755); err != nil { + t.Fatalf("failed to prepare route config dir: %v", err) + } + if err := os.WriteFile( + filepath.Join(filepath.Dir(runner.Config.RouteConfigPath), "atsflare_access.log"), + []byte("{\"ts\":\""+time.Now().UTC().Format(time.RFC3339)+"\",\"host\":\"edge.example.com\",\"remote_addr\":\"10.0.0.8\",\"status\":200}\n"), + 0o644, + ); err != nil { + t.Fatalf("failed to prepare access log: %v", err) + } + + runErr := runner.Run(ctx) + if runErr != context.Canceled { + t.Fatalf("expected run to stop by context cancellation, got %v", runErr) + } + if len(heartbeatService.heartbeatPayloads) != 2 { + t.Fatalf("expected two heartbeat payloads, got %d", len(heartbeatService.heartbeatPayloads)) + } + secondPayload := heartbeatService.heartbeatPayloads[1] + if len(secondPayload.BufferedObservability) != 1 { + t.Fatalf("expected second heartbeat to replay one buffered observation, got %+v", secondPayload.BufferedObservability) + } + + replayable, err := bufferStore.Replayable(0, 0) + if err != nil { + t.Fatalf("Replayable after recovery failed: %v", err) + } + if len(replayable) != 0 { + t.Fatalf("expected buffer to be acked after successful heartbeat, got %+v", replayable) + } +} + func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() diff --git a/atsf_agent/internal/config/config.go b/atsf_agent/internal/config/config.go index 8821f550..75e854d3 100644 --- a/atsf_agent/internal/config/config.go +++ b/atsf_agent/internal/config/config.go @@ -12,12 +12,14 @@ import ( ) const ( - defaultDockerMainConfigRelativePath = "etc/nginx/nginx.conf" - defaultDockerRouteConfigRelativePath = "etc/nginx/conf.d/atsflare_routes.conf" - defaultSupportDirRelativePath = "etc/nginx/support" - defaultDockerStateRelativePath = "var/lib/atsflare/agent-state.json" - defaultDockerOpenRestySupportDir = "/etc/nginx/atsflare-support" - defaultOpenRestyObservabilityPort = 18081 + defaultDockerMainConfigRelativePath = "etc/nginx/nginx.conf" + defaultDockerRouteConfigRelativePath = "etc/nginx/conf.d/atsflare_routes.conf" + defaultSupportDirRelativePath = "etc/nginx/support" + defaultDockerStateRelativePath = "var/lib/atsflare/agent-state.json" + defaultObservabilityBufferRelativePath = "var/lib/atsflare/observability-buffer.json" + defaultDockerOpenRestySupportDir = "/etc/nginx/atsflare-support" + defaultOpenRestyObservabilityPort = 18081 + defaultObservabilityReplayMinutes = 15 ) type Config struct { @@ -38,6 +40,8 @@ type Config struct { SupportDir string `json:"support_dir"` OpenrestySupportDir string `json:"openresty_support_dir"` OpenrestyObservabilityPort int `json:"openresty_observability_port"` + ObservabilityBufferPath string `json:"observability_buffer_path"` + ObservabilityReplayMinutes int `json:"observability_replay_minutes"` StatePath string `json:"state_path"` HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` RequestTimeout MillisecondDuration `json:"request_timeout"` @@ -62,6 +66,8 @@ type configFile struct { LegacyCertDir string `json:"cert_dir"` LegacyOpenrestyCertDir string `json:"openresty_cert_dir"` OpenrestyObservabilityPort int `json:"openresty_observability_port"` + ObservabilityBufferPath string `json:"observability_buffer_path"` + ObservabilityReplayMinutes int `json:"observability_replay_minutes"` StatePath string `json:"state_path"` HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` RequestTimeout MillisecondDuration `json:"request_timeout"` @@ -92,6 +98,8 @@ func Load(path string) (*Config, error) { SupportDir: firstNonEmpty(file.SupportDir, file.LegacyCertDir), OpenrestySupportDir: firstNonEmpty(file.OpenrestySupportDir, file.LegacyOpenrestyCertDir), OpenrestyObservabilityPort: file.OpenrestyObservabilityPort, + ObservabilityBufferPath: file.ObservabilityBufferPath, + ObservabilityReplayMinutes: file.ObservabilityReplayMinutes, StatePath: file.StatePath, HeartbeatInterval: file.HeartbeatInterval, RequestTimeout: file.RequestTimeout, @@ -153,6 +161,12 @@ func applyDefaults(cfg *Config, baseDir string) { if cfg.OpenrestyObservabilityPort <= 0 { cfg.OpenrestyObservabilityPort = defaultOpenRestyObservabilityPort } + if cfg.ObservabilityBufferPath == "" { + cfg.ObservabilityBufferPath = joinManagedPath(cfg.DataDir, defaultObservabilityBufferRelativePath) + } + if cfg.ObservabilityReplayMinutes <= 0 { + cfg.ObservabilityReplayMinutes = defaultObservabilityReplayMinutes + } if cfg.HeartbeatInterval <= 0 { cfg.HeartbeatInterval = MillisecondDuration(10 * time.Second) } @@ -184,6 +198,9 @@ func normalizeManagedPaths(cfg *Config) { if usesSlashPath(cfg.StatePath) { cfg.StatePath = filepath.ToSlash(cfg.StatePath) } + if usesSlashPath(cfg.ObservabilityBufferPath) { + cfg.ObservabilityBufferPath = filepath.ToSlash(cfg.ObservabilityBufferPath) + } } func usesSlashPath(path string) bool { @@ -213,6 +230,9 @@ func validate(cfg *Config) error { if cfg.OpenrestyObservabilityPort <= 0 || cfg.OpenrestyObservabilityPort > 65535 { return errors.New("openresty_observability_port 必须在 1-65535 之间") } + if cfg.ObservabilityReplayMinutes <= 0 { + return errors.New("observability_replay_minutes 必须大于 0") + } return nil } diff --git a/atsf_agent/internal/config/config_test.go b/atsf_agent/internal/config/config_test.go index 143b1124..e5da0747 100644 --- a/atsf_agent/internal/config/config_test.go +++ b/atsf_agent/internal/config/config_test.go @@ -54,9 +54,15 @@ func TestLoadDockerModeUsesManagedPaths(t *testing.T) { if cfg.StatePath != filepath.Join(dir, "data", defaultDockerStateRelativePath) { t.Fatalf("unexpected state path: %s", cfg.StatePath) } + if cfg.ObservabilityBufferPath != filepath.Join(dir, "data", defaultObservabilityBufferRelativePath) { + t.Fatalf("unexpected observability buffer path: %s", cfg.ObservabilityBufferPath) + } if cfg.OpenrestyObservabilityPort != defaultOpenRestyObservabilityPort { t.Fatalf("unexpected openresty observability port: %d", cfg.OpenrestyObservabilityPort) } + if cfg.ObservabilityReplayMinutes != defaultObservabilityReplayMinutes { + t.Fatalf("unexpected observability replay minutes: %d", cfg.ObservabilityReplayMinutes) + } } func TestLoadPathModeKeepsExplicitPaths(t *testing.T) { @@ -94,6 +100,9 @@ func TestLoadPathModeKeepsExplicitPaths(t *testing.T) { if cfg.StatePath != "/tmp/agent-state.json" { t.Fatalf("unexpected state path: %s", cfg.StatePath) } + if cfg.ObservabilityBufferPath != filepath.Join(dir, "data", defaultObservabilityBufferRelativePath) { + t.Fatalf("unexpected observability buffer path: %s", cfg.ObservabilityBufferPath) + } if cfg.OpenrestySupportDir != cfg.SupportDir { t.Fatalf("expected path mode openresty support dir to equal support dir, got %s / %s", cfg.OpenrestySupportDir, cfg.SupportDir) } @@ -134,6 +143,9 @@ func TestLoadUsesCustomDataDirForGeneratedFiles(t *testing.T) { if cfg.StatePath != "/srv/atsflare/"+defaultDockerStateRelativePath { t.Fatalf("unexpected state path: %s", cfg.StatePath) } + if cfg.ObservabilityBufferPath != "/srv/atsflare/"+defaultObservabilityBufferRelativePath { + t.Fatalf("unexpected observability buffer path: %s", cfg.ObservabilityBufferPath) + } if cfg.SupportDir != "/srv/atsflare/"+defaultSupportDirRelativePath { t.Fatalf("unexpected support dir: %s", cfg.SupportDir) } @@ -213,6 +225,9 @@ func TestSavePersistsMillisecondsAndOmitsRuntimeVersions(t *testing.T) { if decoded["openresty_observability_port"] != float64(defaultOpenRestyObservabilityPort) { t.Fatalf("unexpected observability port: %#v", decoded["openresty_observability_port"]) } + if decoded["observability_replay_minutes"] != float64(defaultObservabilityReplayMinutes) { + t.Fatalf("unexpected observability replay minutes: %#v", decoded["observability_replay_minutes"]) + } if _, ok := decoded["nginx_path"]; ok { t.Fatal("legacy nginx_path should not be persisted") } diff --git a/atsf_agent/internal/protocol/agent_api.go b/atsf_agent/internal/protocol/agent_api.go index 996a2727..95ead961 100644 --- a/atsf_agent/internal/protocol/agent_api.go +++ b/atsf_agent/internal/protocol/agent_api.go @@ -36,19 +36,20 @@ 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"` - Profile *NodeSystemProfile `json:"profile,omitempty"` - Snapshot *NodeMetricSnapshot `json:"snapshot,omitempty"` - TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"` - HealthEvents []NodeHealthEvent `json:"health_events"` + 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"` + BufferedObservability []BufferedObservabilityRecord `json:"buffered_observability,omitempty"` + HealthEvents []NodeHealthEvent `json:"health_events"` } type NodeSystemProfile struct { @@ -92,6 +93,12 @@ type NodeTrafficReport struct { SourceCountries map[string]int64 `json:"source_countries"` } +type BufferedObservabilityRecord struct { + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + Snapshot *NodeMetricSnapshot `json:"snapshot,omitempty"` + TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"` +} + type NodeHealthEvent struct { EventType string `json:"event_type"` Severity string `json:"severity"` diff --git a/atsf_agent/internal/state/observability_buffer.go b/atsf_agent/internal/state/observability_buffer.go new file mode 100644 index 00000000..09467fad --- /dev/null +++ b/atsf_agent/internal/state/observability_buffer.go @@ -0,0 +1,171 @@ +package state + +import ( + "encoding/json" + "os" + "path/filepath" + "sort" + "sync" + + "atsflare-agent/internal/protocol" +) + +const observabilityBufferWindowSeconds = 60 + +type ObservabilityBufferRecord struct { + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + Snapshot *protocol.NodeMetricSnapshot `json:"snapshot,omitempty"` + TrafficReport *protocol.NodeTrafficReport `json:"traffic_report,omitempty"` + QueuedAtUnix int64 `json:"queued_at_unix"` +} + +type ObservabilityBufferStore struct { + path string + mu sync.Mutex +} + +func NewObservabilityBufferStore(path string) *ObservabilityBufferStore { + return &ObservabilityBufferStore{path: filepath.Clean(path)} +} + +func (s *ObservabilityBufferStore) Upsert(record ObservabilityBufferRecord, retainAfterUnix int64) error { + if s == nil || record.WindowStartedAtUnix <= 0 || (record.Snapshot == nil && record.TrafficReport == nil) { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + + records, err := s.loadUnlocked() + if err != nil { + return err + } + records = pruneObservabilityBufferRecords(records, retainAfterUnix) + replaced := false + for index := range records { + if records[index].WindowStartedAtUnix != record.WindowStartedAtUnix { + continue + } + records[index] = record + replaced = true + break + } + if !replaced { + records = append(records, record) + } + sort.Slice(records, func(i int, j int) bool { + return records[i].WindowStartedAtUnix < records[j].WindowStartedAtUnix + }) + return s.saveUnlocked(records) +} + +func (s *ObservabilityBufferStore) Replayable(currentWindowStartedAtUnix int64, retainAfterUnix int64) ([]ObservabilityBufferRecord, error) { + if s == nil { + return nil, nil + } + s.mu.Lock() + defer s.mu.Unlock() + + records, err := s.loadUnlocked() + if err != nil { + return nil, err + } + records = pruneObservabilityBufferRecords(records, retainAfterUnix) + if err = s.saveUnlocked(records); err != nil { + return nil, err + } + result := make([]ObservabilityBufferRecord, 0, len(records)) + for _, record := range records { + if currentWindowStartedAtUnix > 0 && record.WindowStartedAtUnix >= currentWindowStartedAtUnix { + continue + } + result = append(result, record) + } + return result, nil +} + +func (s *ObservabilityBufferStore) Ack(windowStartedAtUnix []int64, retainAfterUnix int64) error { + if s == nil || len(windowStartedAtUnix) == 0 { + return nil + } + s.mu.Lock() + defer s.mu.Unlock() + + records, err := s.loadUnlocked() + if err != nil { + return err + } + acked := make(map[int64]struct{}, len(windowStartedAtUnix)) + for _, value := range windowStartedAtUnix { + if value > 0 { + acked[value] = struct{}{} + } + } + filtered := make([]ObservabilityBufferRecord, 0, len(records)) + for _, record := range records { + if _, ok := acked[record.WindowStartedAtUnix]; ok { + continue + } + filtered = append(filtered, record) + } + filtered = pruneObservabilityBufferRecords(filtered, retainAfterUnix) + return s.saveUnlocked(filtered) +} + +func (s *ObservabilityBufferStore) loadUnlocked() ([]ObservabilityBufferRecord, error) { + data, err := os.ReadFile(s.path) + if err != nil { + if os.IsNotExist(err) { + return []ObservabilityBufferRecord{}, nil + } + return nil, err + } + if len(data) == 0 { + return []ObservabilityBufferRecord{}, nil + } + var records []ObservabilityBufferRecord + if err = json.Unmarshal(data, &records); err != nil { + return nil, err + } + return records, nil +} + +func (s *ObservabilityBufferStore) saveUnlocked(records []ObservabilityBufferRecord) error { + if err := os.MkdirAll(filepath.Dir(s.path), 0o755); err != nil { + return err + } + data, err := json.MarshalIndent(records, "", " ") + if err != nil { + return err + } + return os.WriteFile(s.path, data, 0o644) +} + +func ObservabilityWindowStartedAt(snapshot *protocol.NodeMetricSnapshot, traffic *protocol.NodeTrafficReport) int64 { + if traffic != nil && traffic.WindowStartedAtUnix > 0 { + return traffic.WindowStartedAtUnix - (traffic.WindowStartedAtUnix % observabilityBufferWindowSeconds) + } + if snapshot == nil || snapshot.CapturedAtUnix <= 0 { + return 0 + } + return snapshot.CapturedAtUnix - (snapshot.CapturedAtUnix % observabilityBufferWindowSeconds) +} + +func pruneObservabilityBufferRecords(records []ObservabilityBufferRecord, retainAfterUnix int64) []ObservabilityBufferRecord { + if len(records) == 0 { + return []ObservabilityBufferRecord{} + } + filtered := make([]ObservabilityBufferRecord, 0, len(records)) + for _, record := range records { + if record.WindowStartedAtUnix <= 0 { + continue + } + if retainAfterUnix > 0 && record.WindowStartedAtUnix < retainAfterUnix { + continue + } + filtered = append(filtered, record) + } + sort.Slice(filtered, func(i int, j int) bool { + return filtered[i].WindowStartedAtUnix < filtered[j].WindowStartedAtUnix + }) + return filtered +} diff --git a/atsf_agent/internal/state/observability_buffer_test.go b/atsf_agent/internal/state/observability_buffer_test.go new file mode 100644 index 00000000..bc64e49f --- /dev/null +++ b/atsf_agent/internal/state/observability_buffer_test.go @@ -0,0 +1,68 @@ +package state + +import ( + "path/filepath" + "testing" + + "atsflare-agent/internal/protocol" +) + +func TestObservabilityBufferStoreUpsertReplayAndAck(t *testing.T) { + store := NewObservabilityBufferStore(filepath.Join(t.TempDir(), "observability-buffer.json")) + + if err := store.Upsert(ObservabilityBufferRecord{ + WindowStartedAtUnix: 1710403200, + Snapshot: &protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403205}, + TrafficReport: &protocol.NodeTrafficReport{WindowStartedAtUnix: 1710403200, WindowEndedAtUnix: 1710403260, RequestCount: 5}, + QueuedAtUnix: 1710403205, + }, 1710403000); err != nil { + t.Fatalf("first upsert failed: %v", err) + } + if err := store.Upsert(ObservabilityBufferRecord{ + WindowStartedAtUnix: 1710403200, + Snapshot: &protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403255}, + TrafficReport: &protocol.NodeTrafficReport{WindowStartedAtUnix: 1710403200, WindowEndedAtUnix: 1710403260, RequestCount: 12}, + QueuedAtUnix: 1710403255, + }, 1710403000); err != nil { + t.Fatalf("second upsert failed: %v", err) + } + if err := store.Upsert(ObservabilityBufferRecord{ + WindowStartedAtUnix: 1710403260, + Snapshot: &protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403265}, + TrafficReport: &protocol.NodeTrafficReport{WindowStartedAtUnix: 1710403260, WindowEndedAtUnix: 1710403320, RequestCount: 2}, + QueuedAtUnix: 1710403265, + }, 1710403000); err != nil { + t.Fatalf("third upsert failed: %v", err) + } + + records, err := store.Replayable(1710403260, 1710403000) + if err != nil { + t.Fatalf("Replayable failed: %v", err) + } + if len(records) != 1 { + t.Fatalf("expected one replayable record before current window, got %d", len(records)) + } + if records[0].TrafficReport == nil || records[0].TrafficReport.RequestCount != 12 { + t.Fatalf("expected replayable record to keep latest upsert, got %+v", records[0]) + } + + if err = store.Ack([]int64{1710403200}, 1710403000); err != nil { + t.Fatalf("Ack failed: %v", err) + } + records, err = store.Replayable(0, 1710403000) + if err != nil { + t.Fatalf("Replayable after ack failed: %v", err) + } + if len(records) != 1 || records[0].WindowStartedAtUnix != 1710403260 { + t.Fatalf("unexpected records after ack: %+v", records) + } +} + +func TestObservabilityWindowStartedAt(t *testing.T) { + if value := ObservabilityWindowStartedAt(nil, &protocol.NodeTrafficReport{WindowStartedAtUnix: 1710403200}); value != 1710403200 { + t.Fatalf("unexpected traffic window start: %d", value) + } + if value := ObservabilityWindowStartedAt(&protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403259}, nil); value != 1710403200 { + t.Fatalf("unexpected snapshot-derived window start: %d", value) + } +} diff --git a/atsf_server/service/agent.go b/atsf_server/service/agent.go index 12df0937..f743ac3c 100644 --- a/atsf_server/service/agent.go +++ b/atsf_server/service/agent.go @@ -24,19 +24,20 @@ 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"` - Profile *AgentNodeSystemProfile `json:"profile,omitempty"` - Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"` - TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"` - HealthEvents []AgentNodeHealthEvent `json:"health_events"` + 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"` + BufferedObservability []AgentBufferedObservabilityRecord `json:"buffered_observability,omitempty"` + HealthEvents []AgentNodeHealthEvent `json:"health_events"` } type ApplyLogPayload struct { diff --git a/atsf_server/service/node_update_test.go b/atsf_server/service/node_update_test.go index febac672..93d61fe6 100644 --- a/atsf_server/service/node_update_test.go +++ b/atsf_server/service/node_update_test.go @@ -666,6 +666,135 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) { } } +func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) { + setupServiceTestDB(t) + + node := &model.Node{ + NodeID: "node-observe-buffered", + Name: "observe-buffered-edge", + IP: "10.0.0.32", + AgentToken: "token-observe-buffered", + 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) + } + + now := time.Now().UTC() + _, err := HeartbeatNode(node, AgentNodePayload{ + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + AgentVersion: node.AgentVersion, + NginxVersion: node.NginxVersion, + Snapshot: &AgentNodeMetricSnapshot{ + CapturedAtUnix: now.Unix(), + CPUUsagePercent: 25, + MemoryUsedBytes: 2 * 1024 * 1024 * 1024, + MemoryTotalBytes: 8 * 1024 * 1024 * 1024, + }, + TrafficReport: &AgentNodeTrafficReport{ + WindowStartedAtUnix: now.Add(-time.Minute).Unix(), + WindowEndedAtUnix: now.Unix(), + RequestCount: 20, + ErrorCount: 1, + UniqueVisitorCount: 10, + StatusCodes: map[string]int64{"200": 19, "500": 1}, + TopDomains: map[string]int64{"edge.example.com": 20}, + SourceCountries: map[string]int64{"CN": 12}, + }, + BufferedObservability: []AgentBufferedObservabilityRecord{ + { + WindowStartedAtUnix: now.Add(-2 * time.Minute).Unix(), + Snapshot: &AgentNodeMetricSnapshot{ + CapturedAtUnix: now.Add(-2 * time.Minute).Unix(), + CPUUsagePercent: 30, + MemoryUsedBytes: 3 * 1024 * 1024 * 1024, + MemoryTotalBytes: 8 * 1024 * 1024 * 1024, + }, + TrafficReport: &AgentNodeTrafficReport{ + WindowStartedAtUnix: now.Add(-2 * time.Minute).Unix(), + WindowEndedAtUnix: now.Add(-time.Minute).Unix(), + RequestCount: 40, + ErrorCount: 2, + UniqueVisitorCount: 18, + StatusCodes: map[string]int64{"200": 38, "500": 2}, + TopDomains: map[string]int64{"edge.example.com": 40}, + SourceCountries: map[string]int64{"CN": 20}, + }, + }, + }, + }) + if err != nil { + t.Fatalf("expected heartbeat to succeed: %v", err) + } + + 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) != 2 { + t.Fatalf("expected current and buffered snapshots, got %+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) != 2 { + t.Fatalf("expected current and buffered reports, got %+v", reports) + } + + _, err = HeartbeatNode(node, AgentNodePayload{ + NodeID: node.NodeID, + Name: node.Name, + IP: node.IP, + AgentVersion: node.AgentVersion, + NginxVersion: node.NginxVersion, + BufferedObservability: []AgentBufferedObservabilityRecord{ + { + WindowStartedAtUnix: now.Add(-2 * time.Minute).Unix(), + Snapshot: &AgentNodeMetricSnapshot{ + CapturedAtUnix: now.Add(-2 * time.Minute).Unix(), + CPUUsagePercent: 30, + MemoryUsedBytes: 3 * 1024 * 1024 * 1024, + MemoryTotalBytes: 8 * 1024 * 1024 * 1024, + }, + TrafficReport: &AgentNodeTrafficReport{ + WindowStartedAtUnix: now.Add(-2 * time.Minute).Unix(), + WindowEndedAtUnix: now.Add(-time.Minute).Unix(), + RequestCount: 40, + ErrorCount: 2, + UniqueVisitorCount: 18, + StatusCodes: map[string]int64{"200": 38, "500": 2}, + TopDomains: map[string]int64{"edge.example.com": 40}, + SourceCountries: map[string]int64{"CN": 20}, + }, + }, + }, + }) + if err != nil { + t.Fatalf("expected heartbeat replay dedupe to succeed: %v", err) + } + + snapshots, err = model.ListNodeMetricSnapshots(node.NodeID, time.Time{}, 10) + if err != nil { + t.Fatalf("expected node snapshots query to succeed after replay: %v", err) + } + if len(snapshots) != 2 { + t.Fatalf("expected replay dedupe to keep snapshot count stable, got %+v", snapshots) + } + reports, err = model.ListNodeRequestReports(node.NodeID, time.Time{}, 10) + if err != nil { + t.Fatalf("expected node request reports query to succeed after replay: %v", err) + } + if len(reports) != 2 { + t.Fatalf("expected replay dedupe to keep report count stable, got %+v", reports) + } +} + func TestHeartbeatNodeResolvesMissingHealthEvents(t *testing.T) { setupServiceTestDB(t) @@ -754,6 +883,23 @@ func TestGetNodeObservability(t *testing.T) { }); err != nil { t.Fatalf("failed to insert node system profile: %v", err) } + if err := (&model.NodeMetricSnapshot{ + NodeID: node.NodeID, + CapturedAt: time.Now().Add(-time.Hour), + CPUUsagePercent: 60, + MemoryUsedBytes: 14 * 1024 * 1024 * 1024, + MemoryTotalBytes: 16 * 1024 * 1024 * 1024, + StorageUsedBytes: 90 * 1024 * 1024 * 1024, + StorageTotalBytes: 100 * 1024 * 1024 * 1024, + DiskReadBytes: 0, + DiskWriteBytes: 0, + NetworkRxBytes: 2048, + NetworkTxBytes: 4096, + OpenrestyRxBytes: 8192, + OpenrestyTxBytes: 16384, + }).Insert(); err != nil { + t.Fatalf("failed to insert node metric baseline snapshot: %v", err) + } if err := (&model.NodeMetricSnapshot{ NodeID: node.NodeID, CapturedAt: time.Now(), @@ -807,8 +953,11 @@ func TestGetNodeObservability(t *testing.T) { if view.Profile == nil || view.Profile.OSName != "Ubuntu" { t.Fatalf("unexpected profile: %+v", view.Profile) } - if len(view.MetricSnapshots) != 1 { - t.Fatalf("expected 1 metric snapshot, got %d", len(view.MetricSnapshots)) + if len(view.MetricSnapshots) != 2 { + t.Fatalf("expected 2 metric snapshots, got %d", len(view.MetricSnapshots)) + } + if view.MetricSnapshots[0].DiskWriteBytes != 2048 { + t.Fatalf("expected latest metric snapshot to stay intact, got %+v", view.MetricSnapshots[0]) } if len(view.TrafficReports) != 1 || view.TrafficReports[0].RequestCount != 123 { t.Fatalf("unexpected traffic reports: %+v", view.TrafficReports) @@ -923,6 +1072,23 @@ func TestGetDashboardOverview(t *testing.T) { } } + if err := (&model.NodeMetricSnapshot{ + NodeID: "node-dashboard-a", + CapturedAt: now.Add(-time.Hour), + CPUUsagePercent: 40, + MemoryUsedBytes: 4 * 1024 * 1024 * 1024, + MemoryTotalBytes: 8 * 1024 * 1024 * 1024, + StorageUsedBytes: 48 * 1024 * 1024 * 1024, + StorageTotalBytes: 100 * 1024 * 1024 * 1024, + DiskReadBytes: 0, + DiskWriteBytes: 0, + NetworkRxBytes: 100, + NetworkTxBytes: 150, + OpenrestyRxBytes: 300, + OpenrestyTxBytes: 450, + }).Insert(); err != nil { + t.Fatalf("failed to insert node a baseline metric snapshot: %v", err) + } if err := (&model.NodeMetricSnapshot{ NodeID: "node-dashboard-a", CapturedAt: now, @@ -940,6 +1106,23 @@ func TestGetDashboardOverview(t *testing.T) { }).Insert(); err != nil { t.Fatalf("failed to insert node a metric snapshot: %v", err) } + if err := (&model.NodeMetricSnapshot{ + NodeID: "node-dashboard-b", + CapturedAt: now.Add(-time.Hour), + CPUUsagePercent: 88, + MemoryUsedBytes: 14 * 1024 * 1024 * 1024, + MemoryTotalBytes: 16 * 1024 * 1024 * 1024, + StorageUsedBytes: 93 * 1024 * 1024 * 1024, + StorageTotalBytes: 100 * 1024 * 1024 * 1024, + DiskReadBytes: 0, + DiskWriteBytes: 0, + NetworkRxBytes: 200, + NetworkTxBytes: 300, + OpenrestyRxBytes: 500, + OpenrestyTxBytes: 700, + }).Insert(); err != nil { + t.Fatalf("failed to insert node b baseline metric snapshot: %v", err) + } if err := (&model.NodeMetricSnapshot{ NodeID: "node-dashboard-b", CapturedAt: now, diff --git a/atsf_server/service/observability.go b/atsf_server/service/observability.go index 770e9750..73a7618f 100644 --- a/atsf_server/service/observability.go +++ b/atsf_server/service/observability.go @@ -60,6 +60,12 @@ type AgentNodeTrafficReport struct { SourceCountries map[string]int64 `json:"source_countries"` } +type AgentBufferedObservabilityRecord struct { + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"` + TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"` +} + type AgentNodeHealthEvent struct { EventType string `json:"event_type"` Severity string `json:"severity"` @@ -72,7 +78,7 @@ func persistHeartbeatObservability(nodeID string, payload AgentNodePayload, repo if strings.TrimSpace(nodeID) == "" { return } - if payload.Profile == nil && payload.Snapshot == nil && payload.TrafficReport == nil && payload.HealthEvents == nil { + if payload.Profile == nil && payload.Snapshot == nil && payload.TrafficReport == nil && len(payload.BufferedObservability) == 0 && payload.HealthEvents == nil { return } @@ -80,6 +86,9 @@ func persistHeartbeatObservability(nodeID string, payload AgentNodePayload, repo 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 } @@ -97,6 +106,18 @@ func persistHeartbeatObservability(nodeID string, payload AgentNodePayload, repo } } +func persistBufferedObservability(tx *gorm.DB, nodeID string, records []AgentBufferedObservabilityRecord, reportedAt time.Time) error { + for _, record := range records { + if err := persistNodeMetricSnapshot(tx, nodeID, record.Snapshot, reportedAt); err != nil { + return err + } + if err := persistNodeTrafficReport(tx, nodeID, record.TrafficReport, reportedAt); err != nil { + return err + } + } + return nil +} + func persistNodeSystemProfile(tx *gorm.DB, nodeID string, profile *AgentNodeSystemProfile, reportedAt time.Time) error { if profile == nil { return nil @@ -140,7 +161,7 @@ func persistNodeMetricSnapshot(tx *gorm.DB, nodeID string, snapshot *AgentNodeMe OpenrestyConnections: snapshot.OpenrestyConnections, RawJSON: marshalJSON(snapshot), } - return tx.Create(record).Error + return tx.Where("node_id = ? AND captured_at = ?", nodeID, record.CapturedAt).Assign(record).FirstOrCreate(record).Error } func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *AgentNodeTrafficReport, reportedAt time.Time) error { @@ -162,7 +183,7 @@ func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *AgentNodeTraff SourceCountriesJSON: marshalJSON(report.SourceCountries), RawJSON: marshalJSON(report), } - return tx.Create(record).Error + return tx.Where("node_id = ? AND window_started_at = ? AND window_ended_at = ?", nodeID, record.WindowStartedAt, record.WindowEndedAt).Assign(record).FirstOrCreate(record).Error } func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHealthEvent, reportedAt time.Time) error { diff --git a/atsf_server/service/observability_trends.go b/atsf_server/service/observability_trends.go index beed3710..dfd3117f 100644 --- a/atsf_server/service/observability_trends.go +++ b/atsf_server/service/observability_trends.go @@ -2,6 +2,7 @@ package service import ( "atsflare/model" + "sort" "time" ) @@ -148,13 +149,53 @@ func buildDiskIOTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapshot accumulators[index].nodes = make(map[string]struct{}) } + sort.Slice(snapshots, func(i int, j int) bool { + if snapshots[i].CapturedAt.Equal(snapshots[j].CapturedAt) { + return snapshots[i].NodeID < snapshots[j].NodeID + } + return snapshots[i].CapturedAt.Before(snapshots[j].CapturedAt) + }) + + type diskCounterState struct { + read int64 + write int64 + seen bool + } + + previousByNode := make(map[string]diskCounterState, len(snapshots)) + for _, snapshot := range snapshots { + nodeKey := snapshot.NodeID + if nodeKey == "" { + nodeKey = "__unknown__" + } + + previous := previousByNode[nodeKey] + previousByNode[nodeKey] = diskCounterState{ + read: snapshot.DiskReadBytes, + write: snapshot.DiskWriteBytes, + seen: true, + } + if !previous.seen { + continue + } + index, ok := trendBucketIndex(snapshot.CapturedAt, start) if !ok { continue } - points[index].DiskReadBytes += snapshot.DiskReadBytes - points[index].DiskWriteBytes += snapshot.DiskWriteBytes + + readDelta := snapshot.DiskReadBytes - previous.read + writeDelta := snapshot.DiskWriteBytes - previous.write + if readDelta < 0 { + readDelta = 0 + } + if writeDelta < 0 { + writeDelta = 0 + } + + points[index].DiskReadBytes += readDelta + points[index].DiskWriteBytes += writeDelta if snapshot.NodeID != "" { accumulators[index].nodes[snapshot.NodeID] = struct{}{} } diff --git a/atsf_server/service/observability_trends_test.go b/atsf_server/service/observability_trends_test.go new file mode 100644 index 00000000..09727ba6 --- /dev/null +++ b/atsf_server/service/observability_trends_test.go @@ -0,0 +1,32 @@ +package service + +import ( + "atsflare/model" + "testing" + "time" +) + +func TestBuildDiskIOTrendPointsUsesCounterDelta(t *testing.T) { + now := time.Date(2026, 3, 14, 18, 30, 0, 0, time.UTC) + start := trendWindowStart(now) + + points := buildDiskIOTrendPoints(now, []*model.NodeMetricSnapshot{ + { + NodeID: "node-a", + CapturedAt: start.Add(22 * time.Hour), + DiskReadBytes: 100, + DiskWriteBytes: 200, + }, + { + NodeID: "node-a", + CapturedAt: start.Add(23 * time.Hour), + DiskReadBytes: 250, + DiskWriteBytes: 260, + }, + }) + + last := points[len(points)-1] + if last.DiskReadBytes != 150 || last.DiskWriteBytes != 60 { + t.Fatalf("expected disk io trend to use counter delta, got %+v", last) + } +} diff --git a/atsf_server/web/features/dashboard/components/dashboard-overview.tsx b/atsf_server/web/features/dashboard/components/dashboard-overview.tsx index f90895fc..8335e5d7 100644 --- a/atsf_server/web/features/dashboard/components/dashboard-overview.tsx +++ b/atsf_server/web/features/dashboard/components/dashboard-overview.tsx @@ -43,6 +43,13 @@ function formatBytes(value?: number | null) { return `${current.toFixed(current >= 100 || index === 0 ? 0 : 1)} ${units[index]}`; } +function formatBytesPerSecond(value?: number | null, windowSeconds = 1) { + if (!value || value <= 0 || windowSeconds <= 0) { + return '—'; + } + return `${formatBytes(value / windowSeconds)}/s`; +} + function getErrorMessage(error: unknown) { return error instanceof Error ? error.message : '请求失败,请稍后重试。'; } @@ -281,7 +288,7 @@ export function DashboardOverview() { values: overview.trends.network_24h.map( (point) => point.openresty_rx_bytes, ), - valueFormatter: formatBytes, + valueFormatter: (value) => formatBytesPerSecond(value, 3600), }, { label: 'OpenResty 出站', @@ -289,7 +296,7 @@ export function DashboardOverview() { values: overview.trends.network_24h.map( (point) => point.openresty_tx_bytes, ), - valueFormatter: formatBytes, + valueFormatter: (value) => formatBytesPerSecond(value, 3600), }, ]} /> diff --git a/atsf_server/web/features/nodes/components/node-detail-page.tsx b/atsf_server/web/features/nodes/components/node-detail-page.tsx index e6e9b9d5..3e6e2fe4 100644 --- a/atsf_server/web/features/nodes/components/node-detail-page.tsx +++ b/atsf_server/web/features/nodes/components/node-detail-page.tsx @@ -89,6 +89,13 @@ function formatBytes(value?: number | null) { return `${current.toFixed(current >= 100 || index === 0 ? 0 : 1)} ${units[index]}`; } +function formatBytesPerSecond(value?: number | null, windowSeconds = 1) { + if (!value || value <= 0 || windowSeconds <= 0) { + return '—'; + } + return `${formatBytes(value / windowSeconds)}/s`; +} + function formatPercent(value?: number | null) { if (value === undefined || value === null || Number.isNaN(value)) { return '—'; @@ -929,11 +936,17 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) {

入站: - {formatBytes(latestMetricSnapshot.openresty_rx_bytes)} + {formatBytesPerSecond( + latestMetricSnapshot.openresty_rx_bytes, + 60, + )}

出站: - {formatBytes(latestMetricSnapshot.openresty_tx_bytes)} + {formatBytesPerSecond( + latestMetricSnapshot.openresty_tx_bytes, + 60, + )}

@@ -1087,7 +1100,7 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { observability?.trends.network_24h.map( (point) => point.openresty_rx_bytes, ) ?? [], - valueFormatter: formatBytes, + valueFormatter: (value) => formatBytesPerSecond(value, 3600), }, { label: 'OpenResty 出站', @@ -1096,7 +1109,7 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { observability?.trends.network_24h.map( (point) => point.openresty_tx_bytes, ) ?? [], - valueFormatter: formatBytes, + valueFormatter: (value) => formatBytesPerSecond(value, 3600), }, ]} /> diff --git a/docs/app-config.md b/docs/app-config.md index 3c975328..9af9e89a 100644 --- a/docs/app-config.md +++ b/docs/app-config.md @@ -239,6 +239,7 @@ go run ./cmd/agent -config ./agent.json "openresty_container_name": "atsflare-openresty", "openresty_docker_image": "openresty/openresty:alpine", "openresty_observability_port": 18081, + "observability_replay_minutes": 15, "heartbeat_interval": 10000, "request_timeout": 10000 } @@ -259,6 +260,8 @@ go run ./cmd/agent -config ./agent.json "support_dir": "/usr/local/openresty/nginx/conf/support", "openresty_support_dir": "/usr/local/openresty/nginx/conf/support", "openresty_observability_port": 18081, + "observability_buffer_path": "./data/observability-buffer.json", + "observability_replay_minutes": 15, "state_path": "./data/agent-state.json", "heartbeat_interval": 10000, "request_timeout": 10000 @@ -284,7 +287,9 @@ go run ./cmd/agent -config ./agent.json | `route_config_path` | 路由配置文件写入路径 | 否 | 默认为 `data_dir` 下托管路径 | `/etc/nginx/conf.d/atsflare_routes.conf` | | `support_dir` | Agent 在本机写入受管附属文件的目录,当前包含证书与 Lua 观测脚本 | 否 | 默认为 `data_dir` 下托管 support 目录 | `./data/etc/nginx/support` | | `openresty_support_dir` | OpenResty 实际读取受管附属文件的目录 | 否 | 本机模式默认等于 `support_dir`;Docker 模式默认 `/etc/nginx/atsflare-support` | `/usr/local/openresty/nginx/conf/support` | -| `state_path` | Agent 本地状态文件路径 | 否 | 默认为 `data_dir` 下托管状态文件 | `./data/agent-state.json` | +| `observability_buffer_path` | Agent 本地观测补报缓冲文件路径;用于在 server 短暂离线时按时间窗口落盘待补传数据 | 否 | 默认为 `data_dir` 下托管观测缓冲文件 | `./data/var/lib/atsflare/observability-buffer.json` | +| `observability_replay_minutes` | Agent 恢复心跳后允许批量补传的最近观测窗口时长(分钟) | 否 | `15` | `30` | +| `state_path` | Agent 本地状态文件路径 | 否 | 默认为 `data_dir` 下托管状态文件 | `./data/agent-state.json` | | `heartbeat_interval` | 心跳间隔 | 否 | `10000` 毫秒 | `10000` | | `request_timeout` | HTTP 请求超时时间 | 否 | `10000` 毫秒 | `10000` | @@ -297,6 +302,7 @@ go run ./cmd/agent -config ./agent.json * `node_name` 与 `node_ip` 未填写时会自动探测;若自动探测失败,配置校验会报错 * 未配置 `openresty_path` 时,默认为 Docker OpenResty 模式 * `openresty_observability_port` 默认仅绑定本地回环地址;若节点本机已有端口冲突,可改为其他未占用端口 +* `observability_replay_minutes` 只控制“允许补传最近多少分钟的窗口”;超出该窗口的历史观测会在本地自动裁剪 * 为兼容旧节点,Agent 仍可读取历史字段 `cert_dir` / `openresty_cert_dir`,但保存配置时会统一写回 `support_dir` / `openresty_support_dir` * 配置保存时,`agent_version`、`nginx_version` 由程序运行时维护,不需要写入 JSON * 第五版主配置接管完成后,本机模式下应优先通过 `main_config_path` 由 Agent 写入受管主配置,而不是依赖节点手工维护 include 规则 @@ -308,9 +314,10 @@ go run ./cmd/agent -config ./agent.json | 字段 | 默认值 | | --- | --- | | `main_config_path` | 第五版 Docker 模式默认可落在 `data_dir/etc/nginx/nginx.conf`;本机模式建议显式配置 | -| `route_config_path` | `data_dir/etc/nginx/conf.d/atsflare_routes.conf` | +| `route_config_path` | `data_dir/etc/nginx/conf.d/atsflare_routes.conf` | | `support_dir` | `data_dir/etc/nginx/support` | -| `state_path` | `data_dir/var/lib/atsflare/agent-state.json` | +| `observability_buffer_path` | `data_dir/var/lib/atsflare/observability-buffer.json` | +| `state_path` | `data_dir/var/lib/atsflare/agent-state.json` | Docker OpenResty 模式下: @@ -321,6 +328,7 @@ Docker OpenResty 模式下: 补充说明: * Agent 当前会随受管配置一并向 OpenResty 注入 Lua 观测脚本,并在每次 heartbeat 前通过 `http://127.0.0.1:/atsflare/observability` 读取最近窗口请求指标 +* 若 server 短暂离线,Agent 会把最近窗口观测先写入 `observability_buffer_path`,待 heartbeat 恢复后按时间窗口批量补传最近 `observability_replay_minutes` 分钟的数据 * 同一端口还会暴露仅本机可访问的 `stub_status`,用于采集 OpenResty 活动连接数 ### 2.5 Agent 启动示例 diff --git a/docs/deployment.md b/docs/deployment.md index ac532b62..7458e924 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -155,6 +155,7 @@ swag init -g main.go -o docs "openresty_container_name": "atsflare-openresty", "openresty_docker_image": "openresty/openresty:alpine", "openresty_observability_port": 18081, + "observability_replay_minutes": 15, "heartbeat_interval": 10000, "request_timeout": 10000 } @@ -170,6 +171,7 @@ swag init -g main.go -o docs "openresty_container_name": "atsflare-openresty", "openresty_docker_image": "openresty/openresty:alpine", "openresty_observability_port": 18081, + "observability_replay_minutes": 15, "heartbeat_interval": 10000, "request_timeout": 10000 } @@ -185,6 +187,7 @@ swag init -g main.go -o docs * `node_name` 与 `node_ip` 可省略,未填写时自动探测 * 未配置 `openresty_path` 时,默认使用 Docker OpenResty 容器 * Agent 会在受管 OpenResty 中注入 Lua 观测脚本,并通过 `openresty_observability_port` 对本机暴露最近窗口指标与 `stub_status` +* Agent 还会把最近 `observability_replay_minutes` 分钟内未成功上报的观测窗口落盘到本地缓冲文件,待 server 恢复后自动批量补传 ### 3.3 第五版新增部署约束 @@ -193,6 +196,7 @@ swag init -g main.go -o docs * 本机 OpenResty 模式需要为 Agent 显式提供主配置文件写入路径 * Docker OpenResty 模式需要保证主配置、路由配置和证书目录位于同一套受管挂载路径中 * Docker OpenResty 模式会额外挂载一个仅本机可访问的 `127.0.0.1:` 观测端口,用于 Agent 在 heartbeat 前抓取 Lua 窗口指标 +* Docker / 本机两种模式都会在 `data_dir/var/lib/atsflare/observability-buffer.json` 保留最近待补传窗口;若节点磁盘是临时盘,重启后该缓冲也会丢失 * 节点现存手工维护的主配置如继续保留,必须先迁移为 Server 渲染模板的等价配置,再切换到受管模式 * 主配置切换前必须预留回滚副本,并通过一次 `openresty -t` 失败演练验证回滚