fix(openflare): ClickHouse count scan and dashboard traffic metrics

- Scan ClickHouse count()/countIf() aggregates as uint64 before int64 conversion
- Fix access log list count, region stats, and delete pre-count queries
- Aggregate dashboard 24h traffic totals from hourly trend buckets
- Show current hour value in traffic trend chart summary
- Rename docker-compose service to openflare
- Remove redundant config preview section from performance page
This commit is contained in:
ryan
2026-06-19 17:34:19 +08:00
parent 9491b2a744
commit 3187934f72
14 changed files with 281 additions and 88 deletions
+35 -8
View File
@@ -184,11 +184,14 @@ func buildOverviewView(ctx context.Context) (*OverviewView, error) {
}
cpuNodeCount, memoryNodeCount = applyNodeSnapshotMetrics(&nodeHealth, latestSnapshot, view, cpuNodeCount, memoryNodeCount)
applyNodeTrafficMetrics(&nodeHealth, latestTraffic, view)
applyNodeTrafficMetrics(&nodeHealth, latestTraffic)
view.Nodes = append(view.Nodes, nodeHealth)
}
applyTrafficTotalsFromTrend(&view.Traffic, view.Trends.Traffic24h)
applyTrafficRuntimeMetrics(&view.Traffic, latestTrafficReports)
view.Summary.TotalNodes = len(nodes)
if cpuNodeCount > 0 {
view.Capacity.AverageCPUUsagePercent /= float64(cpuNodeCount)
@@ -234,20 +237,44 @@ func applyNodeSnapshotMetrics(nodeHealth *NodeHealth, snapshot *model.OpenFlareM
return cpuNodeCount, memoryNodeCount
}
func applyNodeTrafficMetrics(nodeHealth *NodeHealth, traffic *model.OpenFlareRequestReport, view *OverviewView) {
func applyNodeTrafficMetrics(nodeHealth *NodeHealth, traffic *model.OpenFlareRequestReport) {
if traffic == nil {
return
}
nodeHealth.RequestCount = traffic.RequestCount
nodeHealth.ErrorCount = traffic.ErrorCount
nodeHealth.UniqueVisitorCount = traffic.UniqueVisitorCount
view.Traffic.RequestCount += traffic.RequestCount
view.Traffic.UniqueVisitors += traffic.UniqueVisitorCount
view.Traffic.ErrorCount += traffic.ErrorCount
if duration := traffic.WindowEndedAt.Sub(traffic.WindowStartedAt).Seconds(); duration > 0 {
view.Traffic.EstimatedQPS += float64(traffic.RequestCount) / duration
}
func applyTrafficTotalsFromTrend(traffic *Traffic, points []observability.TrafficTrendPoint) {
if traffic == nil {
return
}
traffic.RequestCount = 0
traffic.ErrorCount = 0
for _, point := range points {
traffic.RequestCount += point.RequestCount
traffic.ErrorCount += point.ErrorCount
}
}
func applyTrafficRuntimeMetrics(traffic *Traffic, latestReports map[string]*model.OpenFlareRequestReport) {
if traffic == nil {
return
}
traffic.UniqueVisitors = 0
traffic.EstimatedQPS = 0
traffic.ReportedNodes = 0
for _, report := range latestReports {
if report == nil {
continue
}
traffic.UniqueVisitors += report.UniqueVisitorCount
traffic.ReportedNodes++
if duration := report.WindowEndedAt.Sub(report.WindowStartedAt).Seconds(); duration > 0 {
traffic.EstimatedQPS += float64(report.RequestCount) / duration
}
}
view.Traffic.ReportedNodes++
}
func compressOverview(view *OverviewView) *OverviewPayload {
@@ -0,0 +1,54 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package observability
import (
"testing"
"time"
"github.com/Rain-kl/Wavelet/internal/model"
)
func TestBuildTrafficTrendPointsBucketsByHour(t *testing.T) {
t.Parallel()
now := time.Date(2026, 6, 19, 17, 30, 0, 0, time.UTC)
reports := []*model.OpenFlareRequestReport{
{
NodeID: "node-a",
WindowStartedAt: now.Add(-3 * time.Hour),
WindowEndedAt: now.Add(-3*time.Hour + time.Minute),
RequestCount: 10,
ErrorCount: 1,
},
{
NodeID: "node-a",
WindowStartedAt: now.Add(-30 * time.Minute),
WindowEndedAt: now.Add(-29 * time.Minute),
RequestCount: 6,
ErrorCount: 0,
},
}
points := BuildTrafficTrendPoints(now, reports)
if len(points) != observabilityTrendBuckets {
t.Fatalf("BuildTrafficTrendPoints() len = %d, want %d", len(points), observabilityTrendBuckets)
}
var totalRequests int64
for _, point := range points {
totalRequests += point.RequestCount
}
if totalRequests != 16 {
t.Fatalf("total request_count = %d, want 16", totalRequests)
}
currentHour := points[len(points)-1]
if currentHour.RequestCount != 6 {
t.Fatalf("current hour request_count = %d, want 6", currentHour.RequestCount)
}
if currentHour.ErrorCount != 0 {
t.Fatalf("current hour error_count = %d, want 0", currentHour.ErrorCount)
}
}
@@ -69,13 +69,6 @@ func ListAccessLogs(ctx context.Context, filter AccessLogFilter, page, pageSize
return logs, safeUint64Count(total), nil
}
func safeUint64Count(count int64) uint64 {
if count < 0 {
return 0
}
return uint64(count)
}
func applyFilter(query *gorm.DB, filter AccessLogFilter) *gorm.DB {
if filter.UserIDs != nil {
if len(filter.UserIDs) == 0 {
@@ -0,0 +1,20 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package analytics
import "math"
func safeUint64Count(count int64) uint64 {
if count < 0 {
return 0
}
return uint64(count)
}
func safeInt64Count(count uint64) int64 {
if count > math.MaxInt64 {
return math.MaxInt64
}
return int64(count)
}
@@ -0,0 +1,54 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package analytics
import (
"math"
"testing"
)
func TestSafeInt64Count(t *testing.T) {
t.Parallel()
tests := []struct {
name string
count uint64
want int64
}{
{name: "zero", count: 0, want: 0},
{name: "small", count: 42, want: 42},
{name: "max int64", count: math.MaxInt64, want: math.MaxInt64},
{name: "overflow clamps", count: math.MaxUint64, want: math.MaxInt64},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
if got := safeInt64Count(tt.count); got != tt.want {
t.Fatalf("safeInt64Count(%d) = %d, want %d", tt.count, got, tt.want)
}
})
}
}
func TestSafeUint64Count(t *testing.T) {
t.Parallel()
tests := []struct {
name string
count int64
want uint64
}{
{name: "zero", count: 0, want: 0},
{name: "positive", count: 42, want: 42},
{name: "negative clamps", count: -1, want: 0},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
if got := safeUint64Count(tt.count); got != tt.want {
t.Fatalf("safeUint64Count(%d) = %d, want %d", tt.count, got, tt.want)
}
})
}
}
@@ -87,7 +87,7 @@ func CountNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) (int64
clause, args := buildNodeAccessLogFilterClause(filter)
tableName := nodeAccessLogTableName()
var totalRecords int64
var totalRecords uint64
countSQL := fmt.Sprintf("SELECT count() FROM %s WHERE %s", tableName, clause)
if err := conn.QueryRow(ctx, countSQL, args...).Scan(&totalRecords); err != nil {
return 0, 0, fmt.Errorf("count node access logs: %w", err)
@@ -100,11 +100,11 @@ SELECT count() FROM (
WHERE %s AND trim(remote_addr) != ''
GROUP BY trimmed_remote_addr
)`, tableName, clause)
var totalIPs int64
var totalIPs uint64
if err := conn.QueryRow(ctx, ipSQL, args...).Scan(&totalIPs); err != nil {
return 0, 0, fmt.Errorf("count node access log ips: %w", err)
}
return totalRecords, totalIPs, nil
return safeInt64Count(totalRecords), safeInt64Count(totalIPs), nil
}
// RegionCountsNodeAccessLogs returns region counts for a node since a time.
@@ -134,11 +134,17 @@ ORDER BY count DESC, trimmed_region ASC`, tableName, clause)
var result []NodeAccessLogRegionCount
for rows.Next() {
var item NodeAccessLogRegionCount
if err := rows.Scan(&item.Region, &item.Count); err != nil {
var (
region string
count uint64
)
if err := rows.Scan(&region, &count); err != nil {
return nil, fmt.Errorf("scan region count row: %w", err)
}
result = append(result, item)
result = append(result, NodeAccessLogRegionCount{
Region: region,
Count: safeInt64Count(count),
})
}
return result, nil
}
@@ -46,7 +46,7 @@ func deleteNodeAccessLogsWithCount(ctx context.Context, countSQL string, countAr
if err != nil {
return 0, err
}
var count int64
var count uint64
if err := conn.QueryRow(ctx, countSQL, countArgs...).Scan(&count); err != nil {
return 0, fmt.Errorf("count node access logs for delete: %w", err)
}
@@ -56,5 +56,5 @@ func deleteNodeAccessLogsWithCount(ctx context.Context, countSQL string, countAr
if err := conn.Exec(ctx, deleteSQL, deleteArgs...); err != nil {
return 0, fmt.Errorf("delete node access logs: %w", err)
}
return count, nil
return safeInt64Count(count), nil
}
@@ -76,11 +76,20 @@ GROUP BY bucket_epoch`, bucketExpr, tableName, clause)
var result []NodeAccessLogBucketAggregate
for rows.Next() {
var item NodeAccessLogBucketAggregate
if err := rows.Scan(&item.BucketEpoch, &item.RequestCount, &item.SuccessCount, &item.ClientErrorCount, &item.ServerErrorCount); err != nil {
var (
bucketEpoch int64
requestCount, successCount, clientErrorCount, serverErrorCount uint64
)
if err := rows.Scan(&bucketEpoch, &requestCount, &successCount, &clientErrorCount, &serverErrorCount); err != nil {
return nil, fmt.Errorf("scan bucket aggregate row: %w", err)
}
result = append(result, item)
result = append(result, NodeAccessLogBucketAggregate{
BucketEpoch: bucketEpoch,
RequestCount: safeInt64Count(requestCount),
SuccessCount: safeInt64Count(successCount),
ClientErrorCount: safeInt64Count(clientErrorCount),
ServerErrorCount: safeInt64Count(serverErrorCount),
})
}
return result, nil
}
@@ -156,11 +165,22 @@ GROUP BY trimmed_remote_addr`, lastSeenExpr, tableName, queryClause)
var result []NodeAccessLogIPAggregate
for rows.Next() {
var item NodeAccessLogIPAggregate
if err := rows.Scan(&item.RemoteAddr, &item.RequestCount, &item.SuccessCount, &item.ClientErrorCount, &item.ServerErrorCount, &item.LastSeenEpoch); err != nil {
var (
remoteAddr string
lastSeenEpoch int64
requestCount, successCount, clientErrorCount, serverErrorCount uint64
)
if err := rows.Scan(&remoteAddr, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &lastSeenEpoch); err != nil {
return nil, fmt.Errorf("scan ip aggregate row: %w", err)
}
result = append(result, item)
result = append(result, NodeAccessLogIPAggregate{
RemoteAddr: remoteAddr,
RequestCount: safeInt64Count(requestCount),
SuccessCount: safeInt64Count(successCount),
ClientErrorCount: safeInt64Count(clientErrorCount),
ServerErrorCount: safeInt64Count(serverErrorCount),
LastSeenEpoch: lastSeenEpoch,
})
}
return result, nil
}
@@ -198,11 +218,20 @@ GROUP BY trimmed_remote_addr`, recentClause, lastSeenExpr, tableName, clause)
var result []NodeAccessLogIPSummary
for rows.Next() {
var item NodeAccessLogIPSummary
if err := rows.Scan(&item.RemoteAddr, &item.TotalRequests, &item.RecentRequests, &item.LastSeenEpoch); err != nil {
var (
remoteAddr string
lastSeenEpoch int64
totalRequests, recentRequests uint64
)
if err := rows.Scan(&remoteAddr, &totalRequests, &recentRequests, &lastSeenEpoch); err != nil {
return nil, fmt.Errorf("scan ip summary row: %w", err)
}
result = append(result, item)
result = append(result, NodeAccessLogIPSummary{
RemoteAddr: remoteAddr,
TotalRequests: safeInt64Count(totalRequests),
RecentRequests: safeInt64Count(recentRequests),
LastSeenEpoch: lastSeenEpoch,
})
}
return result, nil
}
@@ -232,11 +261,17 @@ ORDER BY bucket_epoch ASC`, bucketExpr, tableName, clause)
var result []NodeAccessLogIPTrend
for rows.Next() {
var item NodeAccessLogIPTrend
if err := rows.Scan(&item.BucketEpoch, &item.RequestCount); err != nil {
var (
bucketEpoch int64
requestCount uint64
)
if err := rows.Scan(&bucketEpoch, &requestCount); err != nil {
return nil, fmt.Errorf("scan ip trend row: %w", err)
}
result = append(result, item)
result = append(result, NodeAccessLogIPTrend{
BucketEpoch: bucketEpoch,
RequestCount: safeInt64Count(requestCount),
})
}
return result, nil
}
@@ -109,7 +109,7 @@ func deleteNodeObservabilityWithCount(ctx context.Context, countSQL string, coun
if err != nil {
return 0, err
}
var count int64
var count uint64
if err := conn.QueryRow(ctx, countSQL, countArgs...).Scan(&count); err != nil {
return 0, fmt.Errorf("count node observability rows for delete: %w", err)
}
@@ -119,5 +119,5 @@ func deleteNodeObservabilityWithCount(ctx context.Context, countSQL string, coun
if err := conn.Exec(ctx, deleteSQL, deleteArgs...); err != nil {
return 0, fmt.Errorf("delete node observability rows: %w", err)
}
return count, nil
return safeInt64Count(count), nil
}