mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 14:06:36 +08:00
feat: 添加节点和仪表板的流量分析、趋势和健康状态功能
This commit is contained in:
@@ -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),
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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[];
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -161,9 +161,9 @@
|
||||
|
||||
当前状态:
|
||||
|
||||
* 基本完成
|
||||
* 已落地总览级与节点级聚合接口、24 小时趋势与系统健康摘要
|
||||
* 尚未完成来源国家统计、世界地图相关查询与更完整的告警规则计算
|
||||
* 已完成
|
||||
* 已落地总览级与节点级聚合接口、当前窗口摘要、状态码/域名/来源分布、系统健康摘要与异常节点清单
|
||||
* 已落地最近 24 小时 CPU、内存、网络、磁盘 IO 趋势聚合,并统一总览页与节点详情页的统计口径
|
||||
|
||||
任务:
|
||||
|
||||
|
||||
Reference in New Issue
Block a user