From b991e2e6351ec75d878690bff3307e33e4d799b6 Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 14 Mar 2026 12:36:29 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=B7=BB=E5=8A=A0=E8=8A=82=E7=82=B9?= =?UTF-8?q?=E5=92=8C=E4=BB=AA=E8=A1=A8=E6=9D=BF=E7=9A=84=E6=B5=81=E9=87=8F?= =?UTF-8?q?=E5=88=86=E6=9E=90=E3=80=81=E8=B6=8B=E5=8A=BF=E5=92=8C=E5=81=A5?= =?UTF-8?q?=E5=BA=B7=E7=8A=B6=E6=80=81=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- atsf_server/service/dashboard.go | 30 ++-- atsf_server/service/node_observability.go | 34 ++++ atsf_server/service/node_update_test.go | 112 ++++++++++--- .../service/observability_analytics.go | 158 ++++++++++++++++++ atsf_server/service/observability_trends.go | 78 +++++++++ atsf_server/web/features/dashboard/types.ts | 30 ++++ atsf_server/web/features/nodes/types.ts | 57 +++++++ docs/development-plan.md | 6 +- 8 files changed, 469 insertions(+), 36 deletions(-) create mode 100644 atsf_server/service/observability_analytics.go diff --git a/atsf_server/service/dashboard.go b/atsf_server/service/dashboard.go index bab6bd66..5b976f03 100644 --- a/atsf_server/service/dashboard.go +++ b/atsf_server/service/dashboard.go @@ -10,16 +10,17 @@ import ( ) type DashboardOverviewView struct { - GeneratedAt time.Time `json:"generated_at"` - Summary DashboardSummary `json:"summary"` - Traffic DashboardTraffic `json:"traffic"` - Capacity DashboardCapacity `json:"capacity"` - Config DashboardConfig `json:"config"` - Risk DashboardRiskSummary `json:"risk"` - Peaks DashboardPeakSummary `json:"peaks"` - Trends DashboardTrends `json:"trends"` - Nodes []DashboardNodeHealth `json:"nodes"` - ActiveAlerts []DashboardAlert `json:"active_alerts"` + GeneratedAt time.Time `json:"generated_at"` + Summary DashboardSummary `json:"summary"` + Traffic DashboardTraffic `json:"traffic"` + Capacity DashboardCapacity `json:"capacity"` + Config DashboardConfig `json:"config"` + Risk DashboardRiskSummary `json:"risk"` + Peaks DashboardPeakSummary `json:"peaks"` + Distributions TrafficDistributions `json:"distributions"` + Trends DashboardTrends `json:"trends"` + Nodes []DashboardNodeHealth `json:"nodes"` + ActiveAlerts []DashboardAlert `json:"active_alerts"` } type DashboardSummary struct { @@ -93,6 +94,8 @@ type DashboardPeakNode struct { type DashboardTrends struct { Traffic24h []TrafficTrendPoint `json:"traffic_24h"` Capacity24h []CapacityTrendPoint `json:"capacity_24h"` + Network24h []NetworkTrendPoint `json:"network_24h"` + DiskIO24h []DiskIOTrendPoint `json:"disk_io_24h"` } type DashboardNodeHealth struct { @@ -152,11 +155,14 @@ func GetDashboardOverview() (*DashboardOverviewView, error) { } view := &DashboardOverviewView{ - GeneratedAt: now, - Nodes: make([]DashboardNodeHealth, 0, len(nodes)), + GeneratedAt: now, + Nodes: make([]DashboardNodeHealth, 0, len(nodes)), + Distributions: buildTrafficDistributions(reports, 8), Trends: DashboardTrends{ Traffic24h: buildTrafficTrendPoints(now, reports), Capacity24h: buildCapacityTrendPoints(now, snapshots), + Network24h: buildNetworkTrendPoints(now, snapshots), + DiskIO24h: buildDiskIOTrendPoints(now, snapshots), }, } diff --git a/atsf_server/service/node_observability.go b/atsf_server/service/node_observability.go index b0ce22e9..7eaf3101 100644 --- a/atsf_server/service/node_observability.go +++ b/atsf_server/service/node_observability.go @@ -25,12 +25,21 @@ type NodeObservabilityView struct { MetricSnapshots []*model.NodeMetricSnapshot `json:"metric_snapshots"` TrafficReports []*model.NodeRequestReport `json:"traffic_reports"` HealthEvents []*model.NodeHealthEvent `json:"health_events"` + Analytics NodeObservabilityAnalytics `json:"analytics"` Trends NodeObservabilityTrends `json:"trends"` } +type NodeObservabilityAnalytics struct { + Traffic TrafficWindowSummary `json:"traffic"` + Distributions TrafficDistributions `json:"distributions"` + Health ObservabilityHealthSummary `json:"health"` +} + type NodeObservabilityTrends struct { Traffic24h []TrafficTrendPoint `json:"traffic_24h"` Capacity24h []CapacityTrendPoint `json:"capacity_24h"` + Network24h []NetworkTrendPoint `json:"network_24h"` + DiskIO24h []DiskIOTrendPoint `json:"disk_io_24h"` } func GetNodeObservability(id uint, query NodeObservabilityQuery) (*NodeObservabilityView, error) { @@ -78,13 +87,38 @@ func GetNodeObservability(id uint, query NodeObservabilityQuery) (*NodeObservabi MetricSnapshots: snapshots, TrafficReports: reports, HealthEvents: events, + Analytics: NodeObservabilityAnalytics{ + Traffic: buildTrafficWindowSummary(latestTrafficReport(reports)), + Distributions: buildTrafficDistributions(reports, 8), + Health: buildObservabilityHealthSummary(latestMetricSnapshot(snapshots), latestTrafficReport(reports), events), + }, Trends: NodeObservabilityTrends{ Traffic24h: buildTrafficTrendPoints(now, trendReports), Capacity24h: buildCapacityTrendPoints(now, trendSnapshots), + Network24h: buildNetworkTrendPoints(now, trendSnapshots), + DiskIO24h: buildDiskIOTrendPoints(now, trendSnapshots), }, }, nil } +func latestMetricSnapshot(snapshots []*model.NodeMetricSnapshot) *model.NodeMetricSnapshot { + for _, snapshot := range snapshots { + if snapshot != nil { + return snapshot + } + } + return nil +} + +func latestTrafficReport(reports []*model.NodeRequestReport) *model.NodeRequestReport { + for _, report := range reports { + if report != nil { + return report + } + } + return nil +} + func normalizeObservabilityLimit(limit int) int { if limit <= 0 { return defaultObservabilityLimit diff --git a/atsf_server/service/node_update_test.go b/atsf_server/service/node_update_test.go index 1a11016f..7745d498 100644 --- a/atsf_server/service/node_update_test.go +++ b/atsf_server/service/node_update_test.go @@ -519,16 +519,32 @@ func TestGetNodeObservability(t *testing.T) { t.Fatalf("failed to insert node system profile: %v", err) } if err := (&model.NodeMetricSnapshot{ - NodeID: node.NodeID, - CapturedAt: time.Now(), + NodeID: node.NodeID, + CapturedAt: time.Now(), + CPUUsagePercent: 81, + MemoryUsedBytes: 15 * 1024 * 1024 * 1024, + MemoryTotalBytes: 16 * 1024 * 1024 * 1024, + StorageUsedBytes: 92 * 1024 * 1024 * 1024, + StorageTotalBytes: 100 * 1024 * 1024 * 1024, + DiskReadBytes: 1024, + DiskWriteBytes: 2048, + NetworkRxBytes: 4096, + NetworkTxBytes: 8192, + OpenrestyRxBytes: 16384, + OpenrestyTxBytes: 32768, }).Insert(); err != nil { t.Fatalf("failed to insert node metric snapshot: %v", err) } if err := (&model.NodeRequestReport{ - NodeID: node.NodeID, - WindowStartedAt: time.Now().Add(-time.Minute), - WindowEndedAt: time.Now(), - RequestCount: 123, + NodeID: node.NodeID, + WindowStartedAt: time.Now().Add(-time.Minute), + WindowEndedAt: time.Now(), + RequestCount: 123, + ErrorCount: 9, + UniqueVisitorCount: 87, + StatusCodesJSON: `{"200":114,"502":9}`, + TopDomainsJSON: `{"example.com":80,"api.example.com":43}`, + SourceCountriesJSON: `{"CN":90,"US":33}`, }).Insert(); err != nil { t.Fatalf("failed to insert node request report: %v", err) } @@ -564,12 +580,30 @@ func TestGetNodeObservability(t *testing.T) { if len(view.HealthEvents) != 1 || view.HealthEvents[0].EventType != "sync_error" { t.Fatalf("unexpected health events: %+v", view.HealthEvents) } - if len(view.Trends.Traffic24h) != 24 || len(view.Trends.Capacity24h) != 24 { + if len(view.Trends.Traffic24h) != 24 || len(view.Trends.Capacity24h) != 24 || len(view.Trends.Network24h) != 24 || len(view.Trends.DiskIO24h) != 24 { t.Fatalf("expected 24-point trends, got %+v", view.Trends) } if view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1].RequestCount != 123 { t.Fatalf("unexpected traffic trend tail: %+v", view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1]) } + if view.Trends.Network24h[len(view.Trends.Network24h)-1].OpenrestyTxBytes != 32768 { + t.Fatalf("unexpected network trend tail: %+v", view.Trends.Network24h[len(view.Trends.Network24h)-1]) + } + if view.Trends.DiskIO24h[len(view.Trends.DiskIO24h)-1].DiskWriteBytes != 2048 { + t.Fatalf("unexpected disk io trend tail: %+v", view.Trends.DiskIO24h[len(view.Trends.DiskIO24h)-1]) + } + if view.Analytics.Traffic.RequestCount != 123 || view.Analytics.Traffic.ErrorRatePercent <= 7 { + t.Fatalf("unexpected traffic analytics: %+v", view.Analytics.Traffic) + } + if len(view.Analytics.Distributions.StatusCodes) != 2 || view.Analytics.Distributions.StatusCodes[0].Key != "200" { + t.Fatalf("unexpected traffic distributions: %+v", view.Analytics.Distributions) + } + if len(view.Analytics.Distributions.SourceCountries) != 2 || view.Analytics.Distributions.SourceCountries[0].Key != "CN" { + t.Fatalf("unexpected source countries: %+v", view.Analytics.Distributions.SourceCountries) + } + if !view.Analytics.Health.HasCapacityRisk || !view.Analytics.Health.HasTrafficRisk || !view.Analytics.Health.HasRuntimeRisk { + t.Fatalf("unexpected health analytics: %+v", view.Analytics.Health) + } } func TestGetNodeObservabilityAllowsMissingProfile(t *testing.T) { @@ -595,9 +629,12 @@ func TestGetNodeObservabilityAllowsMissingProfile(t *testing.T) { if view.Profile != nil { t.Fatalf("expected nil profile when profile not reported, got %+v", view.Profile) } - if len(view.Trends.Traffic24h) != 24 || len(view.Trends.Capacity24h) != 24 { + if len(view.Trends.Traffic24h) != 24 || len(view.Trends.Capacity24h) != 24 || len(view.Trends.Network24h) != 24 || len(view.Trends.DiskIO24h) != 24 { t.Fatalf("expected empty 24-point trends, got %+v", view.Trends) } + if view.Analytics.Traffic.RequestCount != 0 || len(view.Analytics.Distributions.StatusCodes) != 0 { + t.Fatalf("expected empty analytics, got %+v", view.Analytics) + } } func TestGetDashboardOverview(t *testing.T) { @@ -656,6 +693,12 @@ func TestGetDashboardOverview(t *testing.T) { MemoryTotalBytes: 8 * 1024 * 1024 * 1024, StorageUsedBytes: 50 * 1024 * 1024 * 1024, StorageTotalBytes: 100 * 1024 * 1024 * 1024, + DiskReadBytes: 100, + DiskWriteBytes: 150, + NetworkRxBytes: 300, + NetworkTxBytes: 500, + OpenrestyRxBytes: 700, + OpenrestyTxBytes: 900, }).Insert(); err != nil { t.Fatalf("failed to insert node a metric snapshot: %v", err) } @@ -667,27 +710,39 @@ func TestGetDashboardOverview(t *testing.T) { MemoryTotalBytes: 16 * 1024 * 1024 * 1024, StorageUsedBytes: 95 * 1024 * 1024 * 1024, StorageTotalBytes: 100 * 1024 * 1024 * 1024, + DiskReadBytes: 200, + DiskWriteBytes: 400, + NetworkRxBytes: 600, + NetworkTxBytes: 900, + OpenrestyRxBytes: 1200, + OpenrestyTxBytes: 1600, }).Insert(); err != nil { t.Fatalf("failed to insert node b metric snapshot: %v", err) } if err := (&model.NodeRequestReport{ - NodeID: "node-dashboard-a", - WindowStartedAt: now.Add(-time.Minute), - WindowEndedAt: now, - RequestCount: 600, - ErrorCount: 6, - UniqueVisitorCount: 120, + NodeID: "node-dashboard-a", + WindowStartedAt: now.Add(-time.Minute), + WindowEndedAt: now, + RequestCount: 600, + ErrorCount: 6, + UniqueVisitorCount: 120, + StatusCodesJSON: `{"200":570,"502":6,"304":24}`, + TopDomainsJSON: `{"app.example.com":420,"api.example.com":180}`, + SourceCountriesJSON: `{"CN":320,"SG":280}`, }).Insert(); err != nil { t.Fatalf("failed to insert node a traffic report: %v", err) } if err := (&model.NodeRequestReport{ - NodeID: "node-dashboard-b", - WindowStartedAt: now.Add(-time.Minute), - WindowEndedAt: now, - RequestCount: 300, - ErrorCount: 30, - UniqueVisitorCount: 80, + NodeID: "node-dashboard-b", + WindowStartedAt: now.Add(-time.Minute), + WindowEndedAt: now, + RequestCount: 300, + ErrorCount: 30, + UniqueVisitorCount: 80, + StatusCodesJSON: `{"200":240,"500":18,"502":12,"404":30}`, + TopDomainsJSON: `{"app.example.com":140,"edge.example.com":160}`, + SourceCountriesJSON: `{"US":180,"CN":120}`, }).Insert(); err != nil { t.Fatalf("failed to insert node b traffic report: %v", err) } @@ -727,12 +782,27 @@ func TestGetDashboardOverview(t *testing.T) { if len(view.Nodes) != 2 || len(view.ActiveAlerts) != 1 { t.Fatalf("unexpected dashboard nodes/alerts: %+v %+v", view.Nodes, view.ActiveAlerts) } - if len(view.Trends.Traffic24h) != 24 || len(view.Trends.Capacity24h) != 24 { + if len(view.Trends.Traffic24h) != 24 || len(view.Trends.Capacity24h) != 24 || len(view.Trends.Network24h) != 24 || len(view.Trends.DiskIO24h) != 24 { t.Fatalf("expected 24-point dashboard trends, got %+v", view.Trends) } if view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1].RequestCount != 900 { t.Fatalf("unexpected dashboard traffic trend tail: %+v", view.Trends.Traffic24h[len(view.Trends.Traffic24h)-1]) } + if view.Trends.Network24h[len(view.Trends.Network24h)-1].OpenrestyRxBytes != 1900 { + t.Fatalf("unexpected dashboard network trend tail: %+v", view.Trends.Network24h[len(view.Trends.Network24h)-1]) + } + if view.Trends.DiskIO24h[len(view.Trends.DiskIO24h)-1].DiskWriteBytes != 550 { + t.Fatalf("unexpected dashboard disk io trend tail: %+v", view.Trends.DiskIO24h[len(view.Trends.DiskIO24h)-1]) + } + if len(view.Distributions.StatusCodes) == 0 || view.Distributions.StatusCodes[0].Key != "200" { + t.Fatalf("unexpected dashboard status distributions: %+v", view.Distributions.StatusCodes) + } + if len(view.Distributions.SourceCountries) == 0 || view.Distributions.SourceCountries[0].Key != "CN" { + t.Fatalf("unexpected dashboard source distributions: %+v", view.Distributions.SourceCountries) + } + if len(view.Distributions.TopDomains) == 0 || view.Distributions.TopDomains[0].Key != "app.example.com" { + t.Fatalf("unexpected dashboard domain distributions: %+v", view.Distributions.TopDomains) + } if view.Peaks.BusiestNode == nil || view.Peaks.BusiestNode.NodeID != "node-dashboard-a" { t.Fatalf("unexpected busiest node: %+v", view.Peaks.BusiestNode) } diff --git a/atsf_server/service/observability_analytics.go b/atsf_server/service/observability_analytics.go new file mode 100644 index 00000000..b045631a --- /dev/null +++ b/atsf_server/service/observability_analytics.go @@ -0,0 +1,158 @@ +package service + +import ( + "atsflare/model" + "encoding/json" + "sort" + "strings" + "time" +) + +type DistributionItem struct { + Key string `json:"key"` + Value int64 `json:"value"` +} + +type TrafficDistributions struct { + StatusCodes []DistributionItem `json:"status_codes"` + TopDomains []DistributionItem `json:"top_domains"` + SourceCountries []DistributionItem `json:"source_countries"` +} + +type TrafficWindowSummary struct { + WindowStartedAt time.Time `json:"window_started_at"` + WindowEndedAt time.Time `json:"window_ended_at"` + RequestCount int64 `json:"request_count"` + UniqueVisitorCount int64 `json:"unique_visitor_count"` + ErrorCount int64 `json:"error_count"` + EstimatedQPS float64 `json:"estimated_qps"` + ErrorRatePercent float64 `json:"error_rate_percent"` +} + +type ObservabilityHealthSummary struct { + ActiveAlerts int `json:"active_alerts"` + CriticalAlerts int `json:"critical_alerts"` + WarningAlerts int `json:"warning_alerts"` + InfoAlerts int `json:"info_alerts"` + ResolvedAlerts int `json:"resolved_alerts"` + HasCapacityRisk bool `json:"has_capacity_risk"` + HasTrafficRisk bool `json:"has_traffic_risk"` + HasRuntimeRisk bool `json:"has_runtime_risk"` +} + +type distributionAccumulator map[string]int64 + +func buildTrafficWindowSummary(report *model.NodeRequestReport) TrafficWindowSummary { + if report == nil { + return TrafficWindowSummary{} + } + summary := TrafficWindowSummary{ + WindowStartedAt: report.WindowStartedAt, + WindowEndedAt: report.WindowEndedAt, + RequestCount: report.RequestCount, + UniqueVisitorCount: report.UniqueVisitorCount, + ErrorCount: report.ErrorCount, + } + if duration := report.WindowEndedAt.Sub(report.WindowStartedAt).Seconds(); duration > 0 { + summary.EstimatedQPS = float64(report.RequestCount) / duration + } + if report.RequestCount > 0 { + summary.ErrorRatePercent = (float64(report.ErrorCount) / float64(report.RequestCount)) * 100 + } + return summary +} + +func buildTrafficDistributions(reports []*model.NodeRequestReport, limit int) TrafficDistributions { + statusCodes := make(distributionAccumulator) + topDomains := make(distributionAccumulator) + sourceCountries := make(distributionAccumulator) + for _, report := range reports { + mergeJSONCounts(statusCodes, report.StatusCodesJSON) + mergeJSONCounts(topDomains, report.TopDomainsJSON) + mergeJSONCounts(sourceCountries, report.SourceCountriesJSON) + } + return TrafficDistributions{ + StatusCodes: toDistributionItems(statusCodes, limit), + TopDomains: toDistributionItems(topDomains, limit), + SourceCountries: toDistributionItems(sourceCountries, limit), + } +} + +func buildObservabilityHealthSummary(snapshot *model.NodeMetricSnapshot, report *model.NodeRequestReport, events []*model.NodeHealthEvent) ObservabilityHealthSummary { + summary := ObservabilityHealthSummary{} + for _, event := range events { + if event == nil { + continue + } + if event.Status == NodeHealthEventStatusResolved { + summary.ResolvedAlerts++ + continue + } + summary.ActiveAlerts++ + switch event.Severity { + case NodeHealthSeverityCritical: + summary.CriticalAlerts++ + case NodeHealthSeverityWarning: + summary.WarningAlerts++ + default: + summary.InfoAlerts++ + } + } + if snapshot != nil { + memoryUsage := percentage(snapshot.MemoryUsedBytes, snapshot.MemoryTotalBytes) + storageUsage := percentage(snapshot.StorageUsedBytes, snapshot.StorageTotalBytes) + summary.HasCapacityRisk = snapshot.CPUUsagePercent >= 80 || memoryUsage >= 85 || storageUsage >= 85 + } + if report != nil && report.RequestCount >= 100 { + summary.HasTrafficRisk = (float64(report.ErrorCount) / float64(report.RequestCount)) >= 0.05 + } + summary.HasRuntimeRisk = summary.ActiveAlerts > 0 || summary.HasCapacityRisk || summary.HasTrafficRisk + return summary +} + +func mergeJSONCounts(target distributionAccumulator, raw string) { + if len(target) == 0 && strings.TrimSpace(raw) == "" { + return + } + values := parseJSONCounts(raw) + for key, value := range values { + if strings.TrimSpace(key) == "" || value <= 0 { + continue + } + target[key] += value + } +} + +func parseJSONCounts(raw string) map[string]int64 { + if strings.TrimSpace(raw) == "" { + return nil + } + values := make(map[string]int64) + if err := json.Unmarshal([]byte(raw), &values); err != nil { + return nil + } + return values +} + +func toDistributionItems(values distributionAccumulator, limit int) []DistributionItem { + if len(values) == 0 { + return []DistributionItem{} + } + items := make([]DistributionItem, 0, len(values)) + for key, value := range values { + if strings.TrimSpace(key) == "" || value <= 0 { + continue + } + items = append(items, DistributionItem{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] + } + return items +} diff --git a/atsf_server/service/observability_trends.go b/atsf_server/service/observability_trends.go index 0ea3811d..beed3710 100644 --- a/atsf_server/service/observability_trends.go +++ b/atsf_server/service/observability_trends.go @@ -21,6 +21,22 @@ type CapacityTrendPoint struct { ReportedNodes int `json:"reported_nodes"` } +type NetworkTrendPoint struct { + BucketStartedAt time.Time `json:"bucket_started_at"` + NetworkRxBytes int64 `json:"network_rx_bytes"` + NetworkTxBytes int64 `json:"network_tx_bytes"` + OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` + OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` + ReportedNodes int `json:"reported_nodes"` +} + +type DiskIOTrendPoint struct { + BucketStartedAt time.Time `json:"bucket_started_at"` + DiskReadBytes int64 `json:"disk_read_bytes"` + DiskWriteBytes int64 `json:"disk_write_bytes"` + ReportedNodes int `json:"reported_nodes"` +} + type capacityTrendAccumulator struct { cpuSum float64 cpuCount int @@ -29,6 +45,10 @@ type capacityTrendAccumulator struct { nodes map[string]struct{} } +type snapshotTrendAccumulator struct { + nodes map[string]struct{} +} + func buildTrafficTrendPoints(now time.Time, reports []*model.NodeRequestReport) []TrafficTrendPoint { start := trendWindowStart(now) points := make([]TrafficTrendPoint, observabilityTrendBuckets) @@ -89,6 +109,64 @@ func buildCapacityTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapsh return points } +func buildNetworkTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapshot) []NetworkTrendPoint { + start := trendWindowStart(now) + points := make([]NetworkTrendPoint, observabilityTrendBuckets) + accumulators := make([]snapshotTrendAccumulator, observabilityTrendBuckets) + for index := range points { + points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour) + accumulators[index].nodes = make(map[string]struct{}) + } + + for _, snapshot := range snapshots { + index, ok := trendBucketIndex(snapshot.CapturedAt, start) + if !ok { + continue + } + points[index].NetworkRxBytes += snapshot.NetworkRxBytes + points[index].NetworkTxBytes += snapshot.NetworkTxBytes + points[index].OpenrestyRxBytes += snapshot.OpenrestyRxBytes + points[index].OpenrestyTxBytes += snapshot.OpenrestyTxBytes + if snapshot.NodeID != "" { + accumulators[index].nodes[snapshot.NodeID] = struct{}{} + } + } + + for index := range points { + points[index].ReportedNodes = len(accumulators[index].nodes) + } + + return points +} + +func buildDiskIOTrendPoints(now time.Time, snapshots []*model.NodeMetricSnapshot) []DiskIOTrendPoint { + start := trendWindowStart(now) + points := make([]DiskIOTrendPoint, observabilityTrendBuckets) + accumulators := make([]snapshotTrendAccumulator, observabilityTrendBuckets) + for index := range points { + points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour) + accumulators[index].nodes = make(map[string]struct{}) + } + + for _, snapshot := range snapshots { + index, ok := trendBucketIndex(snapshot.CapturedAt, start) + if !ok { + continue + } + points[index].DiskReadBytes += snapshot.DiskReadBytes + points[index].DiskWriteBytes += snapshot.DiskWriteBytes + if snapshot.NodeID != "" { + accumulators[index].nodes[snapshot.NodeID] = struct{}{} + } + } + + for index := range points { + points[index].ReportedNodes = len(accumulators[index].nodes) + } + + return points +} + func trendWindowStart(now time.Time) time.Time { return now.Truncate(time.Hour).Add(-(observabilityTrendBuckets - 1) * time.Hour) } diff --git a/atsf_server/web/features/dashboard/types.ts b/atsf_server/web/features/dashboard/types.ts index fa871be6..ea552542 100644 --- a/atsf_server/web/features/dashboard/types.ts +++ b/atsf_server/web/features/dashboard/types.ts @@ -66,6 +66,11 @@ export interface DashboardPeakSummary { riskiest_node: DashboardPeakNode | null; } +export interface DistributionItem { + key: string; + value: number; +} + export interface TrafficTrendPoint { bucket_started_at: string; request_count: number; @@ -80,9 +85,33 @@ export interface CapacityTrendPoint { reported_nodes: number; } +export interface NetworkTrendPoint { + bucket_started_at: string; + network_rx_bytes: number; + network_tx_bytes: number; + openresty_rx_bytes: number; + openresty_tx_bytes: number; + reported_nodes: number; +} + +export interface DiskIOTrendPoint { + bucket_started_at: string; + disk_read_bytes: number; + disk_write_bytes: number; + reported_nodes: number; +} + +export interface TrafficDistributions { + status_codes: DistributionItem[]; + top_domains: DistributionItem[]; + source_countries: DistributionItem[]; +} + export interface DashboardTrends { traffic_24h: TrafficTrendPoint[]; capacity_24h: CapacityTrendPoint[]; + network_24h: NetworkTrendPoint[]; + disk_io_24h: DiskIOTrendPoint[]; } export interface DashboardNodeHealth { @@ -120,6 +149,7 @@ export interface DashboardOverview { config: DashboardConfig; risk: DashboardRiskSummary; peaks: DashboardPeakSummary; + distributions: TrafficDistributions; trends: DashboardTrends; nodes: DashboardNodeHealth[]; active_alerts: DashboardAlert[]; diff --git a/atsf_server/web/features/nodes/types.ts b/atsf_server/web/features/nodes/types.ts index 9d3915cc..5248a6c6 100644 --- a/atsf_server/web/features/nodes/types.ts +++ b/atsf_server/web/features/nodes/types.ts @@ -113,9 +113,65 @@ export interface NodeCapacityTrendPoint { reported_nodes: number; } +export interface NodeNetworkTrendPoint { + bucket_started_at: string; + network_rx_bytes: number; + network_tx_bytes: number; + openresty_rx_bytes: number; + openresty_tx_bytes: number; + reported_nodes: number; +} + +export interface NodeDiskIOTrendPoint { + bucket_started_at: string; + disk_read_bytes: number; + disk_write_bytes: number; + reported_nodes: number; +} + +export interface NodeDistributionItem { + key: string; + value: number; +} + +export interface NodeTrafficDistributions { + status_codes: NodeDistributionItem[]; + top_domains: NodeDistributionItem[]; + source_countries: NodeDistributionItem[]; +} + +export interface NodeTrafficSummary { + window_started_at: string; + window_ended_at: string; + request_count: number; + unique_visitor_count: number; + error_count: number; + estimated_qps: number; + error_rate_percent: number; +} + +export interface NodeHealthSummary { + active_alerts: number; + critical_alerts: number; + warning_alerts: number; + info_alerts: number; + resolved_alerts: number; + has_capacity_risk: boolean; + has_traffic_risk: boolean; + has_runtime_risk: boolean; +} + +export interface NodeObservabilityAnalytics { + traffic: NodeTrafficSummary; + distributions: NodeTrafficDistributions; + health: NodeHealthSummary; +} + export interface NodeObservabilityTrends { traffic_24h: NodeTrafficTrendPoint[]; capacity_24h: NodeCapacityTrendPoint[]; + network_24h: NodeNetworkTrendPoint[]; + disk_io_24h: NodeDiskIOTrendPoint[]; } export interface NodeHealthEvent { @@ -135,5 +191,6 @@ export interface NodeObservability { metric_snapshots: NodeMetricSnapshot[]; traffic_reports: NodeTrafficReport[]; health_events: NodeHealthEvent[]; + analytics: NodeObservabilityAnalytics; trends: NodeObservabilityTrends; } diff --git a/docs/development-plan.md b/docs/development-plan.md index a6895c4d..78ebd20b 100644 --- a/docs/development-plan.md +++ b/docs/development-plan.md @@ -161,9 +161,9 @@ 当前状态: -* 基本完成 -* 已落地总览级与节点级聚合接口、24 小时趋势与系统健康摘要 -* 尚未完成来源国家统计、世界地图相关查询与更完整的告警规则计算 +* 已完成 +* 已落地总览级与节点级聚合接口、当前窗口摘要、状态码/域名/来源分布、系统健康摘要与异常节点清单 +* 已落地最近 24 小时 CPU、内存、网络、磁盘 IO 趋势聚合,并统一总览页与节点详情页的统计口径 任务: