diff --git a/docs/changelog/index.md b/docs/changelog/index.md index c2608d27..da96b340 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -18,6 +18,10 @@ sidebar: false ## [unreleased] +### 修复 + +- 修复节点/仪表盘 24 小时容量、网络、磁盘 IO 趋势在 ClickHouse 限流查询下几乎为空的问题:改为基于 `of_node_metric_snapshots` / `of_node_obs_openresty` 的小时级聚合构建趋势;主机与 OpenResty 累计计数器改为按小时 delta 统计。 + ## [v3.1.1] - 2026-07-06 ### 修改 diff --git a/internal/apps/openflare/dashboard/logics.go b/internal/apps/openflare/dashboard/logics.go index 0a11b8ef..2c1a97d7 100644 --- a/internal/apps/openflare/dashboard/logics.go +++ b/internal/apps/openflare/dashboard/logics.go @@ -137,21 +137,11 @@ func buildOverviewView(ctx context.Context) (*OverviewView, error) { if err != nil { return nil, err } - trafficTrend := observability.BuildTrafficTrendPoints(now, reports) - if trafficHourly, hourlyErr := model.ListOpenFlareTrafficHourlySince(ctx, "", since); hourlyErr == nil && len(trafficHourly) > 0 { - trafficTrend = observability.BuildTrafficTrendPointsFromHourly(now, trafficHourly) - } - view := &OverviewView{ GeneratedAt: now, Nodes: make([]NodeHealth, 0, len(nodes)), Distributions: observability.BuildTrafficDistributions(reports, accessLogRegions, dashboardDistributionLimit), - Trends: observability.NodeTrends{ - Traffic24h: trafficTrend, - Capacity24h: observability.BuildCapacityTrendPoints(now, snapshots), - Network24h: observability.BuildNetworkTrendPoints(now, snapshots, openrestySnapshots), - DiskIO24h: observability.BuildDiskIOTrendPoints(now, snapshots), - }, + Trends: observability.BuildNodeTrends(ctx, now, "", snapshots, openrestySnapshots, reports), } var cpuNodeCount int diff --git a/internal/apps/openflare/observability/analytics.go b/internal/apps/openflare/observability/analytics.go index dce95a8d..b832b80c 100644 --- a/internal/apps/openflare/observability/analytics.go +++ b/internal/apps/openflare/observability/analytics.go @@ -4,6 +4,7 @@ package observability import ( + "context" "encoding/json" "sort" "strings" @@ -13,6 +14,7 @@ import ( ) const observabilityTrendBuckets = 24 +const unknownTrendNodeKey = "__unknown__" const ( healthEventStatusActive = "active" @@ -133,6 +135,12 @@ type diskCounterState struct { seen bool } +type networkCounterState struct { + rx int64 + tx int64 + seen bool +} + func buildTrafficWindowSummary(report *model.OpenFlareRequestReport) *TrafficWindowSummary { if report == nil { return nil @@ -257,6 +265,44 @@ func buildHealthSummary( return summary } +// BuildNodeTrends builds 24h trend series, preferring ClickHouse hourly aggregates +// over limited raw snapshot windows so capacity/network/disk charts stay complete. +func BuildNodeTrends( + ctx context.Context, + now time.Time, + nodeID string, + snapshots []*model.OpenFlareMetricSnapshot, + openrestyObs []*model.OpenFlareNodeObservationOpenresty, + reports []*model.OpenFlareRequestReport, +) NodeTrends { + trendSince := now.Add(-24 * time.Hour) + trafficTrend := BuildTrafficTrendPoints(now, reports) + if trafficHourly, err := model.ListOpenFlareTrafficHourlySince(ctx, nodeID, trendSince); err == nil && len(trafficHourly) > 0 { + trafficTrend = BuildTrafficTrendPointsFromHourly(now, trafficHourly) + } + + capacityTrend := BuildCapacityTrendPoints(now, snapshots) + networkTrend := BuildNetworkTrendPoints(now, snapshots, openrestyObs) + diskIOTrend := BuildDiskIOTrendPoints(now, snapshots) + + metricHourly, metricErr := model.ListOpenFlareMetricHourlySince(ctx, nodeID, trendSince) + if metricErr == nil && len(metricHourly) > 0 { + capacityTrend = BuildCapacityTrendPointsFromHourly(now, metricHourly) + diskIOTrend = BuildDiskIOTrendPointsFromHourly(now, metricHourly) + } + openrestyHourly, openrestyErr := model.ListOpenFlareOpenrestyHourlySince(ctx, nodeID, trendSince) + if metricErr == nil && openrestyErr == nil && (len(metricHourly) > 0 || len(openrestyHourly) > 0) { + networkTrend = BuildNetworkTrendPointsFromHourly(now, metricHourly, openrestyHourly) + } + + return NodeTrends{ + Traffic24h: trafficTrend, + Capacity24h: capacityTrend, + Network24h: networkTrend, + DiskIO24h: diskIOTrend, + } +} + // BuildTrafficTrendPointsFromHourly builds 24h traffic trend buckets from hourly rollups. func BuildTrafficTrendPointsFromHourly(now time.Time, hourly []*model.OpenFlareTrafficHourly) []TrafficTrendPoint { start := trendWindowStart(now) @@ -336,7 +382,30 @@ func BuildCapacityTrendPoints(now time.Time, snapshots []*model.OpenFlareMetricS return points } +// BuildCapacityTrendPointsFromHourly builds 24h capacity trend buckets from hourly aggregates. +func BuildCapacityTrendPointsFromHourly(now time.Time, hourly []*model.OpenFlareMetricHourly) []CapacityTrendPoint { + start := trendWindowStart(now) + points := make([]CapacityTrendPoint, observabilityTrendBuckets) + for index := range points { + points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour) + } + for _, row := range hourly { + if row == nil { + continue + } + index, ok := trendBucketIndex(row.Hour, start) + if !ok { + continue + } + points[index].AverageCPUUsagePercent = row.AverageCPUUsagePercent + points[index].AverageMemoryUsagePercent = row.AverageMemoryUsagePercent + points[index].ReportedNodes = row.ReportedNodes + } + return points +} + // BuildNetworkTrendPoints builds 24h network trend buckets. +// Host and OpenResty counters are cumulative; values are consecutive deltas. func BuildNetworkTrendPoints( now time.Time, snapshots []*model.OpenFlareMetricSnapshot, @@ -349,24 +418,70 @@ func BuildNetworkTrendPoints( points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour) accumulators[index].nodes = make(map[string]struct{}) } + sort.Slice(snapshots, func(i int, j int) bool { + if snapshots[i].CapturedAt.Equal(snapshots[j].CapturedAt) { + return snapshots[i].NodeID < snapshots[j].NodeID + } + return snapshots[i].CapturedAt.Before(snapshots[j].CapturedAt) + }) + previousHostByNode := make(map[string]networkCounterState, len(snapshots)) for _, snapshot := range snapshots { + if snapshot == nil { + continue + } + nodeKey := snapshot.NodeID + if nodeKey == "" { + nodeKey = unknownTrendNodeKey + } + previous := previousHostByNode[nodeKey] + previousHostByNode[nodeKey] = networkCounterState{ + rx: snapshot.NetworkRxBytes, + tx: snapshot.NetworkTxBytes, + seen: true, + } + if !previous.seen { + continue + } index, ok := trendBucketIndex(snapshot.CapturedAt, start) if !ok { continue } - points[index].NetworkRxBytes += snapshot.NetworkRxBytes - points[index].NetworkTxBytes += snapshot.NetworkTxBytes + points[index].NetworkRxBytes += nonNegativeDelta(snapshot.NetworkRxBytes, previous.rx) + points[index].NetworkTxBytes += nonNegativeDelta(snapshot.NetworkTxBytes, previous.tx) if snapshot.NodeID != "" { accumulators[index].nodes[snapshot.NodeID] = struct{}{} } } + sort.Slice(openrestyObs, func(i int, j int) bool { + if openrestyObs[i].CapturedAt.Equal(openrestyObs[j].CapturedAt) { + return openrestyObs[i].NodeID < openrestyObs[j].NodeID + } + return openrestyObs[i].CapturedAt.Before(openrestyObs[j].CapturedAt) + }) + previousOpenrestyByNode := make(map[string]networkCounterState, len(openrestyObs)) for _, obs := range openrestyObs { + if obs == nil { + continue + } + nodeKey := obs.NodeID + if nodeKey == "" { + nodeKey = unknownTrendNodeKey + } + previous := previousOpenrestyByNode[nodeKey] + previousOpenrestyByNode[nodeKey] = networkCounterState{ + rx: obs.OpenrestyRxBytes, + tx: obs.OpenrestyTxBytes, + seen: true, + } + if !previous.seen { + continue + } index, ok := trendBucketIndex(obs.CapturedAt, start) if !ok { continue } - points[index].OpenrestyRxBytes += obs.OpenrestyRxBytes - points[index].OpenrestyTxBytes += obs.OpenrestyTxBytes + points[index].OpenrestyRxBytes += nonNegativeDelta(obs.OpenrestyRxBytes, previous.rx) + points[index].OpenrestyTxBytes += nonNegativeDelta(obs.OpenrestyTxBytes, previous.tx) if obs.NodeID != "" { accumulators[index].nodes[obs.NodeID] = struct{}{} } @@ -377,6 +492,48 @@ func BuildNetworkTrendPoints( return points } +// BuildNetworkTrendPointsFromHourly builds 24h network trend buckets from hourly aggregates. +func BuildNetworkTrendPointsFromHourly( + now time.Time, + metricHourly []*model.OpenFlareMetricHourly, + openrestyHourly []*model.OpenFlareOpenrestyHourly, +) []NetworkTrendPoint { + start := trendWindowStart(now) + points := make([]NetworkTrendPoint, observabilityTrendBuckets) + for index := range points { + points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour) + } + for _, row := range metricHourly { + if row == nil { + continue + } + index, ok := trendBucketIndex(row.Hour, start) + if !ok { + continue + } + points[index].NetworkRxBytes += row.NetworkRxBytes + points[index].NetworkTxBytes += row.NetworkTxBytes + if row.ReportedNodes > points[index].ReportedNodes { + points[index].ReportedNodes = row.ReportedNodes + } + } + for _, row := range openrestyHourly { + if row == nil { + continue + } + index, ok := trendBucketIndex(row.Hour, start) + if !ok { + continue + } + points[index].OpenrestyRxBytes += row.OpenrestyRxBytes + points[index].OpenrestyTxBytes += row.OpenrestyTxBytes + if row.ReportedNodes > points[index].ReportedNodes { + points[index].ReportedNodes = row.ReportedNodes + } + } + return points +} + // BuildDiskIOTrendPoints builds 24h disk IO trend buckets. func BuildDiskIOTrendPoints(now time.Time, snapshots []*model.OpenFlareMetricSnapshot) []DiskIOTrendPoint { start := trendWindowStart(now) @@ -396,7 +553,7 @@ func BuildDiskIOTrendPoints(now time.Time, snapshots []*model.OpenFlareMetricSna for _, snapshot := range snapshots { nodeKey := snapshot.NodeID if nodeKey == "" { - nodeKey = "__unknown__" + nodeKey = unknownTrendNodeKey } previous := previousByNode[nodeKey] previousByNode[nodeKey] = diskCounterState{ @@ -411,16 +568,8 @@ func BuildDiskIOTrendPoints(now time.Time, snapshots []*model.OpenFlareMetricSna if !ok { continue } - readDelta := snapshot.DiskReadBytes - previous.read - writeDelta := snapshot.DiskWriteBytes - previous.write - if readDelta < 0 { - readDelta = 0 - } - if writeDelta < 0 { - writeDelta = 0 - } - points[index].DiskReadBytes += readDelta - points[index].DiskWriteBytes += writeDelta + points[index].DiskReadBytes += nonNegativeDelta(snapshot.DiskReadBytes, previous.read) + points[index].DiskWriteBytes += nonNegativeDelta(snapshot.DiskWriteBytes, previous.write) if snapshot.NodeID != "" { accumulators[index].nodes[snapshot.NodeID] = struct{}{} } @@ -431,6 +580,36 @@ func BuildDiskIOTrendPoints(now time.Time, snapshots []*model.OpenFlareMetricSna return points } +// BuildDiskIOTrendPointsFromHourly builds 24h disk IO trend buckets from hourly aggregates. +func BuildDiskIOTrendPointsFromHourly(now time.Time, hourly []*model.OpenFlareMetricHourly) []DiskIOTrendPoint { + start := trendWindowStart(now) + points := make([]DiskIOTrendPoint, observabilityTrendBuckets) + for index := range points { + points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour) + } + for _, row := range hourly { + if row == nil { + continue + } + index, ok := trendBucketIndex(row.Hour, start) + if !ok { + continue + } + points[index].DiskReadBytes += row.DiskReadBytes + points[index].DiskWriteBytes += row.DiskWriteBytes + points[index].ReportedNodes = row.ReportedNodes + } + return points +} + +func nonNegativeDelta(current int64, previous int64) int64 { + delta := current - previous + if delta < 0 { + return 0 + } + return delta +} + func latestMetricSnapshot(snapshots []*model.OpenFlareMetricSnapshot) *model.OpenFlareMetricSnapshot { var latest *model.OpenFlareMetricSnapshot for _, snapshot := range snapshots { diff --git a/internal/apps/openflare/observability/analytics_test.go b/internal/apps/openflare/observability/analytics_test.go index 967030b2..0ee2535e 100644 --- a/internal/apps/openflare/observability/analytics_test.go +++ b/internal/apps/openflare/observability/analytics_test.go @@ -132,3 +132,84 @@ func TestBuildTrafficWindowSummaryNilWithoutReport(t *testing.T) { t.Fatalf("buildTrafficWindowSummary(nil) = %#v, want nil", summary) } } + +func TestBuildCapacityTrendPointsFromHourlyFillsBuckets(t *testing.T) { + t.Parallel() + + now := time.Date(2026, 7, 10, 9, 30, 0, 0, time.UTC) + hourly := []*model.OpenFlareMetricHourly{ + { + Hour: now.Add(-3 * time.Hour).Truncate(time.Hour), + AverageCPUUsagePercent: 42.5, + AverageMemoryUsagePercent: 61.2, + ReportedNodes: 1, + }, + { + Hour: now.Truncate(time.Hour), + AverageCPUUsagePercent: 12.0, + AverageMemoryUsagePercent: 50.0, + ReportedNodes: 2, + }, + } + + points := BuildCapacityTrendPointsFromHourly(now, hourly) + if len(points) != observabilityTrendBuckets { + t.Fatalf("len = %d, want %d", len(points), observabilityTrendBuckets) + } + if points[len(points)-4].AverageCPUUsagePercent != 42.5 { + t.Fatalf("hour-3 cpu = %v, want 42.5", points[len(points)-4].AverageCPUUsagePercent) + } + if points[len(points)-1].ReportedNodes != 2 { + t.Fatalf("current hour reported_nodes = %d, want 2", points[len(points)-1].ReportedNodes) + } +} + +func TestBuildNetworkTrendPointsUsesCounterDeltas(t *testing.T) { + t.Parallel() + + now := time.Date(2026, 7, 10, 9, 30, 0, 0, time.UTC) + base := now.Truncate(time.Hour) + snapshots := []*model.OpenFlareMetricSnapshot{ + {NodeID: "n1", CapturedAt: base.Add(10 * time.Minute), NetworkRxBytes: 1000, NetworkTxBytes: 2000}, + {NodeID: "n1", CapturedAt: base.Add(20 * time.Minute), NetworkRxBytes: 1500, NetworkTxBytes: 2600}, + } + openrestyObs := []*model.OpenFlareNodeObservationOpenresty{ + {NodeID: "n1", CapturedAt: base.Add(10 * time.Minute), OpenrestyRxBytes: 100, OpenrestyTxBytes: 200}, + {NodeID: "n1", CapturedAt: base.Add(20 * time.Minute), OpenrestyRxBytes: 180, OpenrestyTxBytes: 250}, + } + + points := BuildNetworkTrendPoints(now, snapshots, openrestyObs) + current := points[len(points)-1] + if current.NetworkRxBytes != 500 { + t.Fatalf("network_rx_bytes = %d, want 500", current.NetworkRxBytes) + } + if current.NetworkTxBytes != 600 { + t.Fatalf("network_tx_bytes = %d, want 600", current.NetworkTxBytes) + } + if current.OpenrestyRxBytes != 80 { + t.Fatalf("openresty_rx_bytes = %d, want 80", current.OpenrestyRxBytes) + } + if current.OpenrestyTxBytes != 50 { + t.Fatalf("openresty_tx_bytes = %d, want 50", current.OpenrestyTxBytes) + } +} + +func TestBuildDiskIOTrendPointsFromHourlyFillsBuckets(t *testing.T) { + t.Parallel() + + now := time.Date(2026, 7, 10, 9, 30, 0, 0, time.UTC) + hourly := []*model.OpenFlareMetricHourly{ + { + Hour: now.Add(-1 * time.Hour).Truncate(time.Hour), + DiskReadBytes: 1024, + DiskWriteBytes: 2048, + ReportedNodes: 1, + }, + } + + points := BuildDiskIOTrendPointsFromHourly(now, hourly) + prev := points[len(points)-2] + if prev.DiskReadBytes != 1024 || prev.DiskWriteBytes != 2048 { + t.Fatalf("previous hour disk io = %#v, want read=1024 write=2048", prev) + } +} diff --git a/internal/apps/openflare/observability/node_logics.go b/internal/apps/openflare/observability/node_logics.go index 9efaaa56..1c8e6cad 100644 --- a/internal/apps/openflare/observability/node_logics.go +++ b/internal/apps/openflare/observability/node_logics.go @@ -134,11 +134,6 @@ func GetNodeObservability(ctx context.Context, id uint, query NodeQuery) (*NodeV if err != nil { return nil, err } - trafficTrend := BuildTrafficTrendPoints(now, reports) - if trafficHourly, hourlyErr := model.ListOpenFlareTrafficHourlySince(ctx, node.NodeID, now.Add(-24*time.Hour)); hourlyErr == nil && len(trafficHourly) > 0 { - trafficTrend = BuildTrafficTrendPointsFromHourly(now, trafficHourly) - } - view := &NodeView{ NodeID: node.NodeID, Profile: profile, @@ -150,12 +145,7 @@ func GetNodeObservability(ctx context.Context, id uint, query NodeQuery) (*NodeV Distributions: BuildTrafficDistributions(reports, accessLogRegions, defaultTrafficDistributionLimit), Health: buildHealthSummary(latestMetricSnapshot(snapshots), latestTrafficReport(reports), events), }, - Trends: NodeTrends{ - Traffic24h: trafficTrend, - Capacity24h: BuildCapacityTrendPoints(now, snapshots), - Network24h: BuildNetworkTrendPoints(now, snapshots, openrestyObs), - DiskIO24h: BuildDiskIOTrendPoints(now, snapshots), - }, + Trends: BuildNodeTrends(ctx, now, node.NodeID, snapshots, openrestyObs, reports), } if node.NodeType == "tunnel_relay" { frpsObs, frpsErr := model.ListOpenFlareNodeObservationFrps(ctx, node.NodeID, time.Time{}, 1) diff --git a/internal/model/openflare_observability.go b/internal/model/openflare_observability.go index 5a580ef2..9f769fb5 100644 --- a/internal/model/openflare_observability.go +++ b/internal/model/openflare_observability.go @@ -369,6 +369,72 @@ func ListOpenFlareTrafficHourlySince(ctx context.Context, nodeID string, since t return result, nil } +// OpenFlareMetricHourly is an hourly metric snapshot aggregation row. +type OpenFlareMetricHourly struct { + Hour time.Time `json:"hour"` + AverageCPUUsagePercent float64 `json:"average_cpu_usage_percent"` + AverageMemoryUsagePercent float64 `json:"average_memory_usage_percent"` + NetworkRxBytes int64 `json:"network_rx_bytes"` + NetworkTxBytes int64 `json:"network_tx_bytes"` + DiskReadBytes int64 `json:"disk_read_bytes"` + DiskWriteBytes int64 `json:"disk_write_bytes"` + ReportedNodes int `json:"reported_nodes"` +} + +// OpenFlareOpenrestyHourly is an hourly OpenResty observation aggregation row. +type OpenFlareOpenrestyHourly struct { + Hour time.Time `json:"hour"` + OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` + OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` + ReportedNodes int `json:"reported_nodes"` +} + +// ListOpenFlareMetricHourlySince returns hourly metric aggregates since the given time. +func ListOpenFlareMetricHourlySince(ctx context.Context, nodeID string, since time.Time) ([]*OpenFlareMetricHourly, error) { + rows, err := analyticsrepo.ListNodeMetricHourly(ctx, analyticsrepo.NodeObservabilityFilter{ + NodeID: nodeID, + Since: since, + }) + if err != nil { + return nil, err + } + result := make([]*OpenFlareMetricHourly, len(rows)) + for index, row := range rows { + result[index] = &OpenFlareMetricHourly{ + Hour: row.Hour, + AverageCPUUsagePercent: row.AverageCPUUsagePercent, + AverageMemoryUsagePercent: row.AverageMemoryUsagePercent, + NetworkRxBytes: row.NetworkRxBytes, + NetworkTxBytes: row.NetworkTxBytes, + DiskReadBytes: row.DiskReadBytes, + DiskWriteBytes: row.DiskWriteBytes, + ReportedNodes: row.ReportedNodes, + } + } + return result, nil +} + +// ListOpenFlareOpenrestyHourlySince returns hourly OpenResty aggregates since the given time. +func ListOpenFlareOpenrestyHourlySince(ctx context.Context, nodeID string, since time.Time) ([]*OpenFlareOpenrestyHourly, error) { + rows, err := analyticsrepo.ListNodeOpenrestyHourly(ctx, analyticsrepo.NodeObservabilityFilter{ + NodeID: nodeID, + Since: since, + }) + if err != nil { + return nil, err + } + result := make([]*OpenFlareOpenrestyHourly, len(rows)) + for index, row := range rows { + result[index] = &OpenFlareOpenrestyHourly{ + Hour: row.Hour, + OpenrestyRxBytes: row.OpenrestyRxBytes, + OpenrestyTxBytes: row.OpenrestyTxBytes, + ReportedNodes: row.ReportedNodes, + } + } + return result, nil +} + // ListOpenFlareActiveHealthEvents returns active health events across all nodes. func ListOpenFlareActiveHealthEvents(ctx context.Context) ([]*OpenFlareHealthEvent, error) { conn := db.DB(ctx) diff --git a/internal/repository/analytics/node_observability.go b/internal/repository/analytics/node_observability.go index db1232a8..1b53b4fd 100644 --- a/internal/repository/analytics/node_observability.go +++ b/internal/repository/analytics/node_observability.go @@ -256,6 +256,29 @@ type NodeTrafficHourly struct { UniqueVisitorCount int64 } +// NodeMetricHourly is an hourly metric snapshot aggregation row. +// +// Disk and host network counters are cumulative; deltas use consecutive +// lagInFrame samples per node (negative deltas after counter reset are dropped). +type NodeMetricHourly struct { + Hour time.Time + AverageCPUUsagePercent float64 + AverageMemoryUsagePercent float64 + NetworkRxBytes int64 + NetworkTxBytes int64 + DiskReadBytes int64 + DiskWriteBytes int64 + ReportedNodes int +} + +// NodeOpenrestyHourly is an hourly OpenResty observation aggregation row. +type NodeOpenrestyHourly struct { + Hour time.Time + OpenrestyRxBytes int64 + OpenrestyTxBytes int64 + ReportedNodes int +} + // ListNodeTrafficHourly returns hourly traffic rollup rows matching filter. func ListNodeTrafficHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeTrafficHourly, error) { conn, err := observabilityConn() @@ -298,6 +321,148 @@ ORDER BY hour ASC`, nodeTrafficHourlyTableName, clause) return result, nil } +// ListNodeMetricHourly returns hourly metric snapshot aggregates matching filter. +// Capacity uses sample averages; network/disk use consecutive counter deltas via lagInFrame. +func ListNodeMetricHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeMetricHourly, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "captured_at") + tableName := nodeMetricSnapshotTableName() + sql := fmt.Sprintf(` +SELECT + hour, + avg(cpu_usage_percent) AS average_cpu_usage_percent, + avg(memory_usage_percent) AS average_memory_usage_percent, + sum(if(network_rx_delta >= 0, network_rx_delta, 0)) AS network_rx_bytes, + sum(if(network_tx_delta >= 0, network_tx_delta, 0)) AS network_tx_bytes, + sum(if(disk_read_delta >= 0, disk_read_delta, 0)) AS disk_read_bytes, + sum(if(disk_write_delta >= 0, disk_write_delta, 0)) AS disk_write_bytes, + toUInt64(uniqExact(node_id)) AS reported_nodes +FROM ( + SELECT + node_id, + toStartOfHour(captured_at) AS hour, + cpu_usage_percent, + if(memory_total_bytes > 0, (memory_used_bytes * 100.0) / memory_total_bytes, 0) AS memory_usage_percent, + network_rx_bytes - lagInFrame(network_rx_bytes, 1, network_rx_bytes) OVER ( + PARTITION BY node_id ORDER BY captured_at, id + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW + ) AS network_rx_delta, + network_tx_bytes - lagInFrame(network_tx_bytes, 1, network_tx_bytes) OVER ( + PARTITION BY node_id ORDER BY captured_at, id + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW + ) AS network_tx_delta, + disk_read_bytes - lagInFrame(disk_read_bytes, 1, disk_read_bytes) OVER ( + PARTITION BY node_id ORDER BY captured_at, id + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW + ) AS disk_read_delta, + disk_write_bytes - lagInFrame(disk_write_bytes, 1, disk_write_bytes) OVER ( + PARTITION BY node_id ORDER BY captured_at, id + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW + ) AS disk_write_delta + FROM %s + WHERE %s +) +GROUP BY hour +ORDER BY hour ASC`, tableName, clause) + rows, err := conn.Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list node metric hourly: %w", err) + } + defer func() { _ = rows.Close() }() + + result := make([]NodeMetricHourly, 0) + for rows.Next() { + var ( + item NodeMetricHourly + reportedNodes uint64 + networkRx int64 + networkTx int64 + diskRead int64 + diskWrite int64 + ) + if err := rows.Scan( + &item.Hour, + &item.AverageCPUUsagePercent, + &item.AverageMemoryUsagePercent, + &networkRx, + &networkTx, + &diskRead, + &diskWrite, + &reportedNodes, + ); err != nil { + return nil, fmt.Errorf("scan node metric hourly row: %w", err) + } + item.Hour = item.Hour.UTC() + item.NetworkRxBytes = networkRx + item.NetworkTxBytes = networkTx + item.DiskReadBytes = diskRead + item.DiskWriteBytes = diskWrite + item.ReportedNodes = int(safeInt64Count(reportedNodes)) + result = append(result, item) + } + return result, nil +} + +// ListNodeOpenrestyHourly returns hourly OpenResty observation aggregates matching filter. +func ListNodeOpenrestyHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "captured_at") + tableName := nodeObsOpenrestyTableName() + sql := fmt.Sprintf(` +SELECT + hour, + sum(if(openresty_rx_delta >= 0, openresty_rx_delta, 0)) AS openresty_rx_bytes, + sum(if(openresty_tx_delta >= 0, openresty_tx_delta, 0)) AS openresty_tx_bytes, + toUInt64(uniqExact(node_id)) AS reported_nodes +FROM ( + SELECT + node_id, + toStartOfHour(captured_at) AS hour, + openresty_rx_bytes - lagInFrame(openresty_rx_bytes, 1, openresty_rx_bytes) OVER ( + PARTITION BY node_id ORDER BY captured_at, id + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW + ) AS openresty_rx_delta, + openresty_tx_bytes - lagInFrame(openresty_tx_bytes, 1, openresty_tx_bytes) OVER ( + PARTITION BY node_id ORDER BY captured_at, id + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW + ) AS openresty_tx_delta + FROM %s + WHERE %s +) +GROUP BY hour +ORDER BY hour ASC`, tableName, clause) + rows, err := conn.Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list node openresty hourly: %w", err) + } + defer func() { _ = rows.Close() }() + + result := make([]NodeOpenrestyHourly, 0) + for rows.Next() { + var ( + item NodeOpenrestyHourly + reportedNodes uint64 + rx int64 + tx int64 + ) + if err := rows.Scan(&item.Hour, &rx, &tx, &reportedNodes); err != nil { + return nil, fmt.Errorf("scan node openresty hourly row: %w", err) + } + item.Hour = item.Hour.UTC() + item.OpenrestyRxBytes = rx + item.OpenrestyTxBytes = tx + item.ReportedNodes = int(safeInt64Count(reportedNodes)) + result = append(result, item) + } + return result, nil +} + func scanNodeObsFrpcRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrpc, error) { var result []analyticsmodel.NodeObsFrpc for rows.Next() {