From 3187934f721e5fbd731bcc8123f35f4f96b13e8e Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 19 Jun 2026 17:34:19 +0800 Subject: [PATCH] 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 --- docker-compose.yaml | 4 +- docs/changelog/index.md | 2 + .../dashboard/traffic-trend-chart.tsx | 2 + frontend/app/(main)/performance/page.tsx | 27 +------ frontend/components/data/trend-chart.tsx | 71 +++++++++++++------ internal/apps/openflare/dashboard/logics.go | 43 ++++++++--- .../openflare/observability/analytics_test.go | 54 ++++++++++++++ internal/repository/analytics/access_log.go | 7 -- .../repository/analytics/clickhouse_count.go | 20 ++++++ .../analytics/clickhouse_count_test.go | 54 ++++++++++++++ .../repository/analytics/node_access_log.go | 18 +++-- .../analytics/node_access_log_delete.go | 4 +- .../analytics/node_access_log_stats.go | 59 +++++++++++---- .../analytics/node_observability_delete.go | 4 +- 14 files changed, 281 insertions(+), 88 deletions(-) create mode 100644 internal/apps/openflare/observability/analytics_test.go create mode 100644 internal/repository/analytics/clickhouse_count.go create mode 100644 internal/repository/analytics/clickhouse_count_test.go diff --git a/docker-compose.yaml b/docker-compose.yaml index 6b3961c4..974d4b74 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -1,11 +1,11 @@ services: - wavelet: + openflare: build: context: . dockerfile: docker/Dockerfile args: VERSION: v0.9.9 -# image: ghcr.io/rain-kl/wavelet:latest +# image: ghcr.io/rain-kl/openflare-server:latest restart: unless-stopped env_file: .env environment: diff --git a/docs/changelog/index.md b/docs/changelog/index.md index e115b4bc..7a339289 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -32,7 +32,9 @@ sidebar: false ### 修复 +- 修复总览看板「24 小时请求趋势」摘要误展示 24 小时累计值的问题:趋势图摘要改为「当前小时」桶数据,顶部 24h 统计改为按小时趋势聚合。 - 修复访问日志 ClickHouse 聚合查询因 `trim(x) AS x` 别名与表列同名导致总览看板地域分布及 IP 统计失败的问题。 +- 修复访问日志页 `count()` 扫描类型不匹配(ClickHouse `UInt64` 写入 `int64`)导致列表计数失败的问题。 ### 变更 diff --git a/frontend/app/(main)/components/dashboard/traffic-trend-chart.tsx b/frontend/app/(main)/components/dashboard/traffic-trend-chart.tsx index 19182080..7f02ffce 100644 --- a/frontend/app/(main)/components/dashboard/traffic-trend-chart.tsx +++ b/frontend/app/(main)/components/dashboard/traffic-trend-chart.tsx @@ -30,6 +30,8 @@ export function TrafficTrendChart({ formatTrendHour(point.bucket_started_at))} + summaryScope="last-point" + summaryHint="当前小时" series={[ { label: '请求量', diff --git a/frontend/app/(main)/performance/page.tsx b/frontend/app/(main)/performance/page.tsx index f86d15d3..22d8a1b8 100644 --- a/frontend/app/(main)/performance/page.tsx +++ b/frontend/app/(main)/performance/page.tsx @@ -16,7 +16,7 @@ import {Input} from "@/components/ui/input" import {Label} from "@/components/ui/label" import {Select, SelectContent, SelectItem, SelectTrigger, SelectValue,} from "@/components/ui/select" import {Switch} from "@/components/ui/switch" -import {ConfigVersionService, OptionService} from "@/lib/services/openflare" +import {OptionService} from "@/lib/services/openflare" import { defaultPerformanceFields, @@ -44,12 +44,6 @@ export default function PerformancePage() { enabled: !!user?.is_admin, }) - const previewQuery = useQuery({ - queryKey: ["openflare", "config-preview"], - queryFn: () => ConfigVersionService.preview(), - enabled: !!user?.is_admin, - }) - useEffect(() => { if (!optionsQuery.data) return setFields(mapOptionsToFields(optionsToMap(optionsQuery.data))) @@ -154,25 +148,6 @@ export default function PerformancePage() { -
-
-

发布链路

-

受管模板渲染

-
-
-

预览规则数

-

- {previewQuery.data?.route_count ?? "—"} 条 -

-
-
-

配置预览

-

- 在配置发布页查看完整 nginx 渲染结果 -

-
-
-
diff --git a/frontend/components/data/trend-chart.tsx b/frontend/components/data/trend-chart.tsx index 3da5938e..e05be8e3 100644 --- a/frontend/components/data/trend-chart.tsx +++ b/frontend/components/data/trend-chart.tsx @@ -15,11 +15,16 @@ type TrendChartSeries = { valueFormatter?: (value: number) => string; }; +type TrendChartSummaryScope = 'last-point' | 'total'; + type TrendChartProps = { labels: string[]; series: TrendChartSeries[]; height?: number; yAxisValueFormatter?: (value: number) => string; + showSummary?: boolean; + summaryScope?: TrendChartSummaryScope; + summaryHint?: string; }; type TooltipParam = { @@ -31,11 +36,24 @@ type TooltipParam = { const defaultFormatter = (value: number) => formatCompactNumber(value); +function resolveSummaryValue(values: number[], scope: TrendChartSummaryScope) { + if (values.length === 0) { + return 0; + } + if (scope === 'total') { + return values.reduce((sum, value) => sum + value, 0); + } + return values[values.length - 1] ?? 0; +} + export function TrendChart({ labels, series, height = 220, yAxisValueFormatter, + showSummary = true, + summaryScope = 'last-point', + summaryHint, }: TrendChartProps) { const option = useMemo(() => { const axisFormatter = yAxisValueFormatter ?? defaultFormatter; @@ -165,31 +183,38 @@ export function TrendChart({ return (
-
- {series.map((item) => { - const latestValue = item.values[item.values.length - 1] ?? 0; - const formatter = item.valueFormatter ?? defaultFormatter; - return ( -
-
- -

- {item.label} + {showSummary ? ( +

+ {series.map((item) => { + const summaryValue = resolveSummaryValue(item.values, summaryScope); + const formatter = item.valueFormatter ?? defaultFormatter; + return ( +
+
+ +

+ {item.label} +

+
+

+ {formatter(summaryValue)}

+ {summaryHint ? ( +

+ {summaryHint} +

+ ) : null}
-

- {formatter(latestValue)} -

-
- ); - })} -
+ ); + })} +
+ ) : null}
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 { diff --git a/internal/apps/openflare/observability/analytics_test.go b/internal/apps/openflare/observability/analytics_test.go new file mode 100644 index 00000000..41d20f86 --- /dev/null +++ b/internal/apps/openflare/observability/analytics_test.go @@ -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) + } +} \ No newline at end of file diff --git a/internal/repository/analytics/access_log.go b/internal/repository/analytics/access_log.go index bb210d8e..317000cf 100644 --- a/internal/repository/analytics/access_log.go +++ b/internal/repository/analytics/access_log.go @@ -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 { diff --git a/internal/repository/analytics/clickhouse_count.go b/internal/repository/analytics/clickhouse_count.go new file mode 100644 index 00000000..1b19b475 --- /dev/null +++ b/internal/repository/analytics/clickhouse_count.go @@ -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) +} \ No newline at end of file diff --git a/internal/repository/analytics/clickhouse_count_test.go b/internal/repository/analytics/clickhouse_count_test.go new file mode 100644 index 00000000..c440cd73 --- /dev/null +++ b/internal/repository/analytics/clickhouse_count_test.go @@ -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) + } + }) + } +} \ No newline at end of file diff --git a/internal/repository/analytics/node_access_log.go b/internal/repository/analytics/node_access_log.go index 441433de..c5ca2d2c 100644 --- a/internal/repository/analytics/node_access_log.go +++ b/internal/repository/analytics/node_access_log.go @@ -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(®ion, &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 } diff --git a/internal/repository/analytics/node_access_log_delete.go b/internal/repository/analytics/node_access_log_delete.go index 40c3e7a8..a0593cc8 100644 --- a/internal/repository/analytics/node_access_log_delete.go +++ b/internal/repository/analytics/node_access_log_delete.go @@ -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 } \ No newline at end of file diff --git a/internal/repository/analytics/node_access_log_stats.go b/internal/repository/analytics/node_access_log_stats.go index fbdb9d14..34d79d8e 100644 --- a/internal/repository/analytics/node_access_log_stats.go +++ b/internal/repository/analytics/node_access_log_stats.go @@ -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 } diff --git a/internal/repository/analytics/node_observability_delete.go b/internal/repository/analytics/node_observability_delete.go index e772bf81..1f7c0fa3 100644 --- a/internal/repository/analytics/node_observability_delete.go +++ b/internal/repository/analytics/node_observability_delete.go @@ -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 }