mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-06 07:36:37 +08:00
fix(obs): 对齐无兼容层与健康/UV 权威语义
Agent 本地旧观测缓冲直接删除并运行重建;设计文档去掉兼容期表述。 健康当前态以 PG status/message 为准,CH 仅存 status 与连接时序; Zone 曲线标明分桶 UV,顶部为整窗独立访客。
This commit is contained in:
@@ -3,10 +3,12 @@ package state
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/agent/protocol"
|
||||
@@ -15,6 +17,8 @@ import (
|
||||
const observabilityBufferWindowSeconds = 60
|
||||
|
||||
// ObservabilityBufferRecord stores observability facts for a single time window.
|
||||
// Disk JSON is schema-v2 only: host_metrics / edge_health / access_logs.
|
||||
// Pre-v2 buffers are discarded on load (binary upgrade without data-dir wipe).
|
||||
type ObservabilityBufferRecord struct {
|
||||
WindowStartedAtUnix int64 `json:"window_started_at_unix"`
|
||||
HostMetrics *protocol.NodeMetricSnapshot `json:"host_metrics,omitempty"`
|
||||
@@ -191,10 +195,14 @@ func (s *ObservabilityBufferStore) loadUnlocked() ([]ObservabilityBufferRecord,
|
||||
s.cacheLoaded = true
|
||||
return []ObservabilityBufferRecord{}, nil
|
||||
}
|
||||
var records []ObservabilityBufferRecord
|
||||
if err = json.Unmarshal(data, &records); err != nil {
|
||||
return nil, err
|
||||
|
||||
// Binary upgrade: drop pre-v2 or corrupt buffer entirely; agent rebuilds on subsequent heartbeats.
|
||||
records, reason, ok := parseObservabilityBufferDisk(data)
|
||||
if !ok {
|
||||
s.discardBufferFile(reason)
|
||||
return []ObservabilityBufferRecord{}, nil
|
||||
}
|
||||
|
||||
s.cache = records
|
||||
s.cacheLoaded = true
|
||||
copied := make([]ObservabilityBufferRecord, len(s.cache))
|
||||
@@ -202,6 +210,38 @@ func (s *ObservabilityBufferStore) loadUnlocked() ([]ObservabilityBufferRecord,
|
||||
return copied, nil
|
||||
}
|
||||
|
||||
// parseObservabilityBufferDisk returns v2 records, or ok=false when the on-disk file should be wiped.
|
||||
func parseObservabilityBufferDisk(data []byte) (records []ObservabilityBufferRecord, reason string, ok bool) {
|
||||
raw := strings.TrimSpace(string(data))
|
||||
if raw == "" {
|
||||
return []ObservabilityBufferRecord{}, "", true
|
||||
}
|
||||
// Valid buffer is a JSON array of window records.
|
||||
if !strings.HasPrefix(raw, "[") {
|
||||
return nil, "legacy or unreadable observability buffer", false
|
||||
}
|
||||
// Pre-v2 keys: discard whole file (no field migration).
|
||||
if strings.Contains(raw, `"snapshot"`) ||
|
||||
strings.Contains(raw, `"openresty_observation"`) ||
|
||||
strings.Contains(raw, `"traffic_report"`) {
|
||||
return nil, "legacy observability buffer format", false
|
||||
}
|
||||
if err := json.Unmarshal(data, &records); err != nil {
|
||||
return nil, "observability buffer JSON decode failed", false
|
||||
}
|
||||
return records, "", true
|
||||
}
|
||||
|
||||
func (s *ObservabilityBufferStore) discardBufferFile(reason string) {
|
||||
if err := os.Remove(s.path); err != nil && !os.IsNotExist(err) {
|
||||
slog.Warn("remove observability buffer failed", "path", s.path, "reason", reason, "error", err)
|
||||
} else {
|
||||
slog.Info("discarded observability buffer; will rebuild on run", "path", s.path, "reason", reason)
|
||||
}
|
||||
s.cache = []ObservabilityBufferRecord{}
|
||||
s.cacheLoaded = true
|
||||
}
|
||||
|
||||
func (s *ObservabilityBufferStore) saveUnlocked(records []ObservabilityBufferRecord) error {
|
||||
if err := os.MkdirAll(filepath.Dir(s.path), stateDirPerm); err != nil {
|
||||
return err
|
||||
@@ -210,7 +250,7 @@ func (s *ObservabilityBufferStore) saveUnlocked(records []ObservabilityBufferRec
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := os.WriteFile(s.path, data, stateFilePerm); err != nil {
|
||||
if err := os.WriteFile(s.path, data, stateFilePerm); err != nil { //nolint:gosec // path is agent-local buffer path from config
|
||||
return err
|
||||
}
|
||||
s.cache = records
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
package state
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/agent/protocol"
|
||||
@@ -95,3 +97,93 @@ func TestObservabilityWindowStartedAt(t *testing.T) {
|
||||
t.Fatalf("unexpected host-metrics window start: %d", value)
|
||||
}
|
||||
}
|
||||
|
||||
func TestObservabilityBufferStoreDiscardsLegacyDiskJSON(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "observability-buffer.json")
|
||||
legacy := `[{
|
||||
"window_started_at_unix": 1710403200,
|
||||
"snapshot": {"captured_at_unix": 1710403205, "cpu_usage_percent": 11.5},
|
||||
"openresty_observation": {"captured_at_unix": 1710403206, "openresty_connections": 7},
|
||||
"traffic_report": {"request_count": 42},
|
||||
"access_logs": [{"logged_at_unix": 1710403201, "path": "/", "status_code": 200}]
|
||||
}]`
|
||||
if err := os.WriteFile(path, []byte(legacy), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile: %v", err)
|
||||
}
|
||||
|
||||
store := NewObservabilityBufferStore(path)
|
||||
records, err := store.Replayable(0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("Replayable: %v", err)
|
||||
}
|
||||
if len(records) != 0 {
|
||||
t.Fatalf("expected legacy buffer discarded, got %+v", records)
|
||||
}
|
||||
// Replayable may rewrite an empty v2 array; legacy keys must be gone.
|
||||
if body, err := os.ReadFile(path); err == nil {
|
||||
raw := string(body)
|
||||
for _, key := range []string{`"snapshot"`, `"openresty_observation"`, `"traffic_report"`} {
|
||||
if strings.Contains(raw, key) {
|
||||
t.Fatalf("legacy key %s still present: %s", key, raw)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Fresh upsert after discard should create a clean v2 file.
|
||||
if err := store.Upsert(ObservabilityBufferRecord{
|
||||
WindowStartedAtUnix: 1710403200,
|
||||
HostMetrics: &protocol.NodeMetricSnapshot{CapturedAtUnix: 1710403205},
|
||||
}, 0); err != nil {
|
||||
t.Fatalf("Upsert after discard: %v", err)
|
||||
}
|
||||
records, err = store.Replayable(0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("Replayable after rebuild: %v", err)
|
||||
}
|
||||
if len(records) != 1 || records[0].HostMetrics == nil {
|
||||
t.Fatalf("expected rebuilt buffer, got %+v", records)
|
||||
}
|
||||
}
|
||||
|
||||
func TestObservabilityBufferStoreDiscardsCorruptJSON(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "observability-buffer.json")
|
||||
if err := os.WriteFile(path, []byte(`{not-json`), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile: %v", err)
|
||||
}
|
||||
store := NewObservabilityBufferStore(path)
|
||||
records, err := store.Replayable(0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("Replayable should not fail: %v", err)
|
||||
}
|
||||
if len(records) != 0 {
|
||||
t.Fatalf("expected empty after discard, got %+v", records)
|
||||
}
|
||||
// Corrupt payload must not remain; empty rewrite is fine.
|
||||
if body, err := os.ReadFile(path); err == nil && strings.Contains(string(body), "not-json") {
|
||||
t.Fatalf("corrupt content still on disk: %s", body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestObservabilityBufferStoreKeepsModernJSON(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "observability-buffer.json")
|
||||
modern := `[{
|
||||
"window_started_at_unix": 1710403200,
|
||||
"host_metrics": {"captured_at_unix": 1710403205, "cpu_usage_percent": 3},
|
||||
"edge_health": {"captured_at_unix": 1710403205, "status": "healthy", "connections": 2},
|
||||
"access_logs": []
|
||||
}]`
|
||||
if err := os.WriteFile(path, []byte(modern), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile: %v", err)
|
||||
}
|
||||
store := NewObservabilityBufferStore(path)
|
||||
records, err := store.Replayable(0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("Replayable: %v", err)
|
||||
}
|
||||
if len(records) != 1 || records[0].HostMetrics == nil || records[0].HostMetrics.CPUUsagePercent != 3 {
|
||||
t.Fatalf("modern buffer should be kept: %+v", records)
|
||||
}
|
||||
if _, err := os.Stat(path); err != nil {
|
||||
t.Fatalf("modern buffer file should remain: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -64,6 +64,20 @@ func normalizeNodePayload(payload NodePayload) NodePayload {
|
||||
payload.LastError = truncateForDatabase(payload.LastError, maxDatabaseTextLength)
|
||||
payload.OpenrestyStatus = normalizeOpenrestyStatus(payload.OpenrestyStatus)
|
||||
payload.OpenrestyMessage = truncateForDatabase(payload.OpenrestyMessage, maxDatabaseTextLength)
|
||||
// Align L2 edge_health with top-level status/message (PG is latest-state authority).
|
||||
if payload.EdgeHealth != nil {
|
||||
if s := strings.TrimSpace(payload.EdgeHealth.Status); s != "" {
|
||||
if payload.OpenrestyStatus == "" || payload.OpenrestyStatus == openrestyStatusUnknown {
|
||||
payload.OpenrestyStatus = normalizeOpenrestyStatus(s)
|
||||
}
|
||||
}
|
||||
if m := strings.TrimSpace(payload.EdgeHealth.Message); m != "" && payload.OpenrestyMessage == "" {
|
||||
payload.OpenrestyMessage = truncateForDatabase(m, maxDatabaseTextLength)
|
||||
}
|
||||
// CH series status must match the same authority as PG after normalize.
|
||||
payload.EdgeHealth.Status = payload.OpenrestyStatus
|
||||
payload.EdgeHealth.Message = payload.OpenrestyMessage
|
||||
}
|
||||
return payload
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user