mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-08 00:26:37 +08:00
feat: add observability buffer for heartbeat data
- Introduced `ObservabilityBufferStore` to manage buffered observability records. - Implemented methods for upserting, replaying, and acknowledging records in the buffer. - Enhanced `AgentNodePayload` and `NodePayload` to include buffered observability data. - Updated tests to cover new functionality for observability buffering. - Modified existing services to persist and handle buffered observability data during heartbeats. - Added configuration options for observability buffer path and replay minutes. - Updated documentation to reflect new observability features and configurations.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user