feat: add 24-hour traffic and capacity trend charts to dashboard and node detail pages

- Implemented TrendChart component for visualizing traffic and capacity trends.
- Enhanced DashboardOverview and NodeDetailPage components to display 24-hour request and capacity trends.
- Updated types to include traffic and capacity trend data structures.
- Created observability trends service to aggregate traffic and capacity data.
- Added tests for traffic report generation and trend data handling.
- Introduced new dependencies for charting (echarts and echarts-for-react).
This commit is contained in:
ryan
2026-03-14 11:48:05 +08:00
parent e94132a7d9
commit 16bb9e5879
25 changed files with 1063 additions and 109 deletions
+2
View File
@@ -313,6 +313,7 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload {
}
profile := observability.BuildProfile(r.Config, r.StateStore)
metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore)
trafficReport := observability.BuildTrafficReport(r.Config, r.StateStore)
healthEvents := observability.BuildHealthEvents(snapshot)
return protocol.NodePayload{
NodeID: nodeID,
@@ -326,6 +327,7 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload {
OpenrestyMessage: snapshot.OpenrestyMessage,
Profile: profile,
Snapshot: metricSnapshot,
TrafficReport: trafficReport,
HealthEvents: healthEvents,
}
}
+20 -2
View File
@@ -282,7 +282,8 @@ func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) {
}
func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) {
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
tempDir := t.TempDir()
stateStore := state.NewStore(filepath.Join(tempDir, "state.json"))
if err := stateStore.Save(&state.Snapshot{
NodeID: "node-observe",
CurrentVersion: "20260314-001",
@@ -299,11 +300,22 @@ func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) {
NodeIP: "10.0.0.51",
AgentVersion: config.AgentVersion,
NginxVersion: "1.27.1.2",
DataDir: t.TempDir(),
DataDir: tempDir,
RouteConfigPath: filepath.Join(tempDir, "conf.d", "atsflare_routes.conf"),
HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond),
},
StateStore: stateStore,
}
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\":\"2026-03-14T10:00:00Z\",\"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)
}
firstPayload := runner.nodePayload("node-observe")
if firstPayload.Profile == nil {
@@ -312,6 +324,9 @@ func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) {
if firstPayload.Snapshot == nil {
t.Fatal("expected first heartbeat payload to include metric snapshot")
}
if firstPayload.TrafficReport == nil || firstPayload.TrafficReport.RequestCount != 1 {
t.Fatalf("expected first heartbeat payload to include traffic report, got %+v", firstPayload.TrafficReport)
}
if len(firstPayload.HealthEvents) != 2 {
t.Fatalf("expected health events for openresty and sync error, got %+v", firstPayload.HealthEvents)
}
@@ -323,6 +338,9 @@ func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) {
if secondPayload.Snapshot == nil {
t.Fatal("expected metric snapshot to continue reporting on subsequent heartbeat")
}
if secondPayload.TrafficReport != nil {
t.Fatalf("expected unchanged traffic window to be omitted on subsequent heartbeat, got %+v", secondPayload.TrafficReport)
}
}
func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) {
+20 -4
View File
@@ -19,8 +19,10 @@ import (
const CertDirPlaceholder = "__ATSF_CERT_DIR__"
const RouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__"
const AccessLogPlaceholder = "__ATSF_ACCESS_LOG__"
const DockerMainConfigPath = "/usr/local/openresty/nginx/conf/nginx.conf"
const DockerRouteConfigPath = "/etc/nginx/conf.d/atsflare_routes.conf"
const DockerAccessLogPath = "/etc/nginx/conf.d/atsflare_access.log"
const dockerRuntimeCommand = "openresty"
@@ -285,6 +287,9 @@ func (m *Manager) CurrentChecksum() (string, error) {
if includePath := m.routeConfigIncludePath(); includePath != "" {
normalizedMain = strings.ReplaceAll(normalizedMain, includePath, RouteConfigPlaceholder)
}
if accessLogPath := m.accessLogRuntimePath(); accessLogPath != "" {
normalizedMain = strings.ReplaceAll(normalizedMain, accessLogPath, AccessLogPlaceholder)
}
normalizedRoute := string(data)
if m.NginxCertDir != "" {
normalizedRoute = strings.ReplaceAll(normalizedRoute, m.NginxCertDir, CertDirPlaceholder)
@@ -609,11 +614,14 @@ func (m *Manager) renderRouteConfig(content string) string {
}
func (m *Manager) renderMainConfig(content string) string {
includePath := m.routeConfigIncludePath()
if includePath == "" {
return content
rendered := content
if includePath := m.routeConfigIncludePath(); includePath != "" {
rendered = strings.ReplaceAll(rendered, RouteConfigPlaceholder, includePath)
}
return strings.ReplaceAll(content, RouteConfigPlaceholder, includePath)
if accessLogPath := m.accessLogRuntimePath(); accessLogPath != "" {
rendered = strings.ReplaceAll(rendered, AccessLogPlaceholder, accessLogPath)
}
return rendered
}
func (m *Manager) routeConfigIncludePath() string {
@@ -623,6 +631,14 @@ func (m *Manager) routeConfigIncludePath() string {
return strings.TrimSpace(m.RouteConfigPath)
}
func (m *Manager) accessLogRuntimePath() string {
includePath := m.routeConfigIncludePath()
if strings.TrimSpace(includePath) == "" {
return ""
}
return filepath.ToSlash(filepath.Join(filepath.Dir(includePath), "atsflare_access.log"))
}
func checksum(content string) string {
sum := sha256.Sum256([]byte(content))
return hex.EncodeToString(sum[:])
+8 -6
View File
@@ -337,7 +337,7 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) {
err := manager.Apply(
context.Background(),
"include __ATSF_ROUTE_CONFIG__;\n",
"include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n",
"ssl_certificate __ATSF_CERT_DIR__/1.crt;\n",
[]protocol.SupportFile{{Path: "1.crt", Content: "cert"}},
)
@@ -349,7 +349,8 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) {
if err != nil {
t.Fatalf("failed to read main config: %v", err)
}
if string(mainData) != "include "+routePath+";\n" {
expectedMain := "include " + routePath + ";\naccess_log " + filepath.Join(filepath.Dir(routePath), "atsflare_access.log") + " atsflare_json;\n"
if string(mainData) != expectedMain {
t.Fatalf("unexpected main config: %s", string(mainData))
}
@@ -366,7 +367,7 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) {
t.Fatalf("CurrentChecksum failed: %v", err)
}
expected := bundleChecksum(
"include __ATSF_ROUTE_CONFIG__;\n",
"include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n",
"ssl_certificate __ATSF_CERT_DIR__/1.crt;\n",
[]protocol.SupportFile{{Path: "1.crt", Content: "cert"}},
)
@@ -388,7 +389,7 @@ func TestManagerApplyUsesRuntimeRouteConfigPath(t *testing.T) {
Executor: &fakeExecutor{},
}
if err := manager.Apply(context.Background(), "include __ATSF_ROUTE_CONFIG__;\n", "server { listen 80; }\n", nil); err != nil {
if err := manager.Apply(context.Background(), "include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n", "server { listen 80; }\n", nil); err != nil {
t.Fatalf("Apply failed: %v", err)
}
@@ -396,7 +397,8 @@ func TestManagerApplyUsesRuntimeRouteConfigPath(t *testing.T) {
if err != nil {
t.Fatalf("failed to read main config: %v", err)
}
if string(mainData) != "include "+DockerRouteConfigPath+";\n" {
expectedMain := "include " + DockerRouteConfigPath + ";\naccess_log " + DockerAccessLogPath + " atsflare_json;\n"
if string(mainData) != expectedMain {
t.Fatalf("unexpected main config include path: %s", string(mainData))
}
@@ -405,7 +407,7 @@ func TestManagerApplyUsesRuntimeRouteConfigPath(t *testing.T) {
t.Fatalf("CurrentChecksum failed: %v", err)
}
expected := bundleChecksum(
"include __ATSF_ROUTE_CONFIG__;\n",
"include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n",
"server { listen 80; }\n",
nil,
)
@@ -0,0 +1,202 @@
package observability
import (
"atsflare-agent/internal/config"
"atsflare-agent/internal/protocol"
"atsflare-agent/internal/state"
"bufio"
"encoding/json"
"errors"
"io"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"time"
)
type accessLogRecord struct {
Timestamp string `json:"ts"`
Host string `json:"host"`
RemoteAddr string `json:"remote_addr"`
Status int `json:"status"`
}
type trafficAggregate struct {
windowStartedAt time.Time
windowEndedAt time.Time
requestCount int64
errorCount int64
statusCodes map[string]int64
topDomains map[string]int64
visitors map[string]struct{}
}
func BuildTrafficReport(cfg *config.Config, stateStore *state.Store) *protocol.NodeTrafficReport {
if cfg == nil || stateStore == nil {
return nil
}
snapshot, err := stateStore.Load()
if err != nil {
return nil
}
logPath := managedAccessLogPath(cfg)
file, err := os.Open(logPath)
if err != nil {
if os.IsNotExist(err) {
if snapshot.AccessLogOffset != 0 {
snapshot.AccessLogOffset = 0
_ = stateStore.Save(snapshot)
}
return nil
}
return nil
}
defer file.Close()
info, err := file.Stat()
if err != nil {
return nil
}
offset := snapshot.AccessLogOffset
if offset < 0 || offset > info.Size() {
offset = 0
}
if _, err = file.Seek(offset, io.SeekStart); err != nil {
return nil
}
reader := bufio.NewReader(file)
currentOffset := offset
aggregate := newTrafficAggregate()
for {
line, readErr := reader.ReadBytes('\n')
if len(line) > 0 {
currentOffset += int64(len(line))
aggregate.consume(line)
}
if errors.Is(readErr, io.EOF) {
break
}
if readErr != nil {
return nil
}
}
snapshot.AccessLogOffset = currentOffset
_ = stateStore.Save(snapshot)
return aggregate.report()
}
func managedAccessLogPath(cfg *config.Config) string {
if cfg == nil || strings.TrimSpace(cfg.RouteConfigPath) == "" {
return ""
}
return filepath.Join(filepath.Dir(cfg.RouteConfigPath), "atsflare_access.log")
}
func newTrafficAggregate() *trafficAggregate {
return &trafficAggregate{
statusCodes: make(map[string]int64),
topDomains: make(map[string]int64),
visitors: make(map[string]struct{}),
}
}
func (aggregate *trafficAggregate) consume(line []byte) {
trimmed := strings.TrimSpace(string(line))
if trimmed == "" {
return
}
var record accessLogRecord
if err := json.Unmarshal([]byte(trimmed), &record); err != nil {
return
}
timestamp, err := parseAccessLogTime(record.Timestamp)
if err != nil {
return
}
if aggregate.windowStartedAt.IsZero() || timestamp.Before(aggregate.windowStartedAt) {
aggregate.windowStartedAt = timestamp
}
if aggregate.windowEndedAt.IsZero() || timestamp.After(aggregate.windowEndedAt) {
aggregate.windowEndedAt = timestamp
}
aggregate.requestCount++
if record.Status >= 500 {
aggregate.errorCount++
}
if record.Status > 0 {
aggregate.statusCodes[strconv.Itoa(record.Status)]++
}
if host := strings.TrimSpace(record.Host); host != "" {
aggregate.topDomains[host]++
}
if remoteAddr := strings.TrimSpace(record.RemoteAddr); remoteAddr != "" {
aggregate.visitors[remoteAddr] = struct{}{}
}
}
func (aggregate *trafficAggregate) report() *protocol.NodeTrafficReport {
if aggregate.requestCount == 0 || aggregate.windowStartedAt.IsZero() || aggregate.windowEndedAt.IsZero() {
return nil
}
return &protocol.NodeTrafficReport{
WindowStartedAtUnix: aggregate.windowStartedAt.Unix(),
WindowEndedAtUnix: aggregate.windowEndedAt.Unix(),
RequestCount: aggregate.requestCount,
ErrorCount: aggregate.errorCount,
UniqueVisitorCount: int64(len(aggregate.visitors)),
StatusCodes: cloneTrafficCounts(aggregate.statusCodes, 0),
TopDomains: topCounts(aggregate.topDomains, 8),
SourceCountries: map[string]int64{},
}
}
func parseAccessLogTime(value string) (time.Time, error) {
return time.Parse(time.RFC3339, strings.TrimSpace(value))
}
func cloneTrafficCounts(values map[string]int64, limit int) map[string]int64 {
if len(values) == 0 {
return map[string]int64{}
}
items := make([]trafficCountItem, 0, len(values))
for key, value := range values {
items = append(items, trafficCountItem{key: key, value: value})
}
sort.Slice(items, func(i int, j int) bool {
if items[i].value == items[j].value {
return items[i].key < items[j].key
}
return items[i].value > items[j].value
})
if limit > 0 && len(items) > limit {
items = items[:limit]
}
result := make(map[string]int64, len(items))
for _, item := range items {
result[item.key] = item.value
}
return result
}
type trafficCountItem struct {
key string
value int64
}
func topCounts(values map[string]int64, limit int) map[string]int64 {
return cloneTrafficCounts(values, limit)
}
@@ -0,0 +1,77 @@
package observability
import (
"os"
"path/filepath"
"testing"
"atsflare-agent/internal/config"
"atsflare-agent/internal/state"
)
func TestBuildTrafficReportAggregatesManagedAccessLog(t *testing.T) {
tempDir := t.TempDir()
routeConfigPath := filepath.Join(tempDir, "conf.d", "atsflare_routes.conf")
if err := os.MkdirAll(filepath.Dir(routeConfigPath), 0o755); err != nil {
t.Fatalf("MkdirAll failed: %v", err)
}
logPath := filepath.Join(filepath.Dir(routeConfigPath), "atsflare_access.log")
content := []byte(
"{\"ts\":\"2026-03-14T08:00:00Z\",\"host\":\"app.example.com\",\"remote_addr\":\"10.0.0.1\",\"status\":200}\n" +
"{\"ts\":\"2026-03-14T08:00:05Z\",\"host\":\"app.example.com\",\"remote_addr\":\"10.0.0.2\",\"status\":503}\n" +
"{\"ts\":\"2026-03-14T08:00:08Z\",\"host\":\"api.example.com\",\"remote_addr\":\"10.0.0.1\",\"status\":200}\n",
)
if err := os.WriteFile(logPath, content, 0o644); err != nil {
t.Fatalf("WriteFile failed: %v", err)
}
stateStore := state.NewStore(filepath.Join(tempDir, "state.json"))
report := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore)
if report == nil {
t.Fatal("expected traffic report")
}
if report.RequestCount != 3 || report.ErrorCount != 1 || report.UniqueVisitorCount != 2 {
t.Fatalf("unexpected traffic report counters: %+v", report)
}
if report.StatusCodes["200"] != 2 || report.StatusCodes["503"] != 1 {
t.Fatalf("unexpected status codes: %+v", report.StatusCodes)
}
if report.TopDomains["app.example.com"] != 2 || report.TopDomains["api.example.com"] != 1 {
t.Fatalf("unexpected top domains: %+v", report.TopDomains)
}
snapshot, err := stateStore.Load()
if err != nil {
t.Fatalf("Load failed: %v", err)
}
if snapshot.AccessLogOffset != int64(len(content)) {
t.Fatalf("unexpected access log offset: %d", snapshot.AccessLogOffset)
}
secondReport := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore)
if secondReport != nil {
t.Fatalf("expected no report without appended lines, got %+v", secondReport)
}
}
func TestBuildTrafficReportResetsOffsetAfterTruncate(t *testing.T) {
tempDir := t.TempDir()
routeConfigPath := filepath.Join(tempDir, "conf.d", "atsflare_routes.conf")
if err := os.MkdirAll(filepath.Dir(routeConfigPath), 0o755); err != nil {
t.Fatalf("MkdirAll failed: %v", err)
}
logPath := filepath.Join(filepath.Dir(routeConfigPath), "atsflare_access.log")
if err := os.WriteFile(logPath, []byte("{\"ts\":\"2026-03-14T09:00:00Z\",\"host\":\"app.example.com\",\"remote_addr\":\"10.0.0.3\",\"status\":200}\n"), 0o644); err != nil {
t.Fatalf("WriteFile failed: %v", err)
}
stateStore := state.NewStore(filepath.Join(tempDir, "state.json"))
if err := stateStore.Save(&state.Snapshot{AccessLogOffset: 4096}); err != nil {
t.Fatalf("Save failed: %v", err)
}
report := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore)
if report == nil || report.RequestCount != 1 {
t.Fatalf("expected one request after truncate reset, got %+v", report)
}
}
+1
View File
@@ -20,6 +20,7 @@ type Snapshot struct {
LastCPUStatTotal uint64 `json:"last_cpu_stat_total"`
LastCPUStatIdle uint64 `json:"last_cpu_stat_idle"`
LastMetricAtUnix int64 `json:"last_metric_at_unix"`
AccessLogOffset int64 `json:"access_log_offset"`
}
type Store struct {